Haystack 接入 Langfuse 全流程指南:用 LangfuseConnector 追踪 LLM 管道运行
【免费下载链接】haystackOpen-source AI orchestration framework for building context-engineered, production-ready LLM applications. Design modular pipelines and agent workflows with explicit control over retrieval, routing, memory, and generation. Built for scalable agents, RAG, multimodal applications, semantic search, and conversational systems.项目地址: https://gitcode.com/GitHub_Trending/ha/haystack
Langfuse 是 LLM 应用的可观测性平台,而 Haystack 通过LangfuseConnector这个"零连线"组件,把管道内每一次运行的完整链路——提示词、模型调用、组件输入输出、Token 消耗与耗时——自动送入 Langfuse 仪表盘。读完本文,你将掌握 LangfuseConnector 的安装配置、环境变量语义、管道与 Agent 两种接入方式、flush 数据刷新策略,以及通过SpanHandler深度定制追踪内容的高级玩法。
为什么要在 Haystack 管道中接入 Langfuse
LangfuseConnector的作用是把 Haystack LLM 框架与 Langfuse 连接起来,实现对管道内各组件操作与数据流的追踪(tracing)。它的使用方式非常特殊:把它加进管道,但不需要与任何其他组件连线——当追踪被启用时,它会自动捕获管道内的所有运行操作。LangfuseConnector 的使用价值集中在三个方面(见 LangfuseConnector 组件文档):
- 监控模型性能:查看每次运行的 Token 用量与成本;
- 定位改进空间:识别低质量输出、收集用户反馈,据此优化管道;
- 沉淀数据集:从管道执行结果中构造用于微调和测试的数据集。
它的工作方式是:把LangfuseConnector加入管道后正常运行管道,然后到 Langfuse 网站上查看追踪数据。由于它不参与任何数据流,它只是在管道后台"伴跑",因此它可以放在管道中的任意位置。
核心设计:一个无需连线的"伴跑"追踪组件
要理解 LangfuseConnector 为何无需连线即可生效,需要先了解 Haystack 的追踪架构。在 haystack/tracing/tracer.py 中,Haystack 维护了一个全局tracer(ProxyTracer实例),并通过enable_tracing()/disable_tracing()在空实现NullTracer与真实追踪器之间切换。管道运行时内部大量调用tracing.tracer.trace(...)来创建 span:
- 在 haystack/core/pipeline/pipeline.py#L404-L412 中,每次
pipeline.run()都会开启一个名为haystack.pipeline.run的 span,并附带管道输入数据、元数据等 tags; - 在 haystack/core/pipeline/base.py#L1097-L1124 中,每个组件执行前都会创建
haystack.component.run子 span,并打上haystack.component.name、haystack.component.type、输入输出规格等标签。
LangfuseConnector本质上是在管道启动时把自己注册为全局 tracer(对应LangfuseTracer,见下文"内部机制"一节),于是上述所有 span 都会被自动转发到 Langfuse,这就是它"自动追踪所有管道操作"的原理。当管道运行完成后,run()方法会返回一个包含name、trace_url、trace_id三个键的字典,方便你把追踪链接打印出来或接入手工审计流程。
环境准备与安装
使用LangfuseConnector前需要准备以下内容:
- 拥有一个可用的 Langfuse 账号(Langfuse 云服务或自托管实例均可);
- 在账号设置中获取公钥与私钥;
- 安装集成包:
pip install langfuse-haystack安装完成后,需要配置如下环境变量:
| 环境变量 | 是否必填 | 说明 |
|---|---|---|
LANGFUSE_SECRET_KEY | 必填 | Langfuse API 私钥,对应__init__中的secret_key参数 |
LANGFUSE_PUBLIC_KEY | 必填 | Langfuse API 公钥,对应__init__中的public_key参数 |
HAYSTACK_CONTENT_TRACING_ENABLED | 必填(须为"true") | 开启内容追踪,让 span 携带组件的输入/输出等具体内容 |
LANGFUSE_HOST | 可选 | Langfuse API 地址,默认https://cloud.langfuse.com |
HAYSTACK_LANGFUSE_ENFORCE_FLUSH | 可选 | 设为"false"可关闭每个组件执行后的强制 flush(详见下文 flush 一节) |
重要:必须先设置环境变量,再导入任何 Haystack 组件。这是因为 Haystack 在导入过程中就会初始化内部的追踪组件(参见 haystack/tracing/tracer.py#L118 中
is_content_tracing_enabled的读取逻辑)。更稳妥的做法是在运行脚本前于 shell 中设置这些变量,把配置与代码分离,便于不同环境的管理。
快速上手:在 RAG 对话管道中接入 Langfuse 追踪
下面是一个最小可用示例:用ChatPromptBuilder+OpenAIChatGenerator构建一条对话管道,并把LangfuseConnector作为名为tracer的组件加入。每次运行会生成一条 trace,包含完整执行上下文(提示词、模型回复、元数据等),输出中的 URL 就是 Langfuse 中该次运行的追踪链接。
import os os.environ["HAYSTACK_CONTENT_TRACING_ENABLED"] = "true" from haystack import Pipeline from haystack.components.builders import ChatPromptBuilder from haystack.components.generators.chat import OpenAIChatGenerator from haystack.dataclasses import ChatMessage from haystack_integrations.components.connectors.langfuse import ( LangfuseConnector, ) pipe = Pipeline() pipe.add_component("tracer", LangfuseConnector("Chat example")) pipe.add_component("prompt_builder", ChatPromptBuilder()) pipe.add_component("llm", OpenAIChatGenerator(model="gpt-4o-mini")) pipe.connect("prompt_builder.prompt", "llm.messages") messages = [ ChatMessage.from_system( "Always respond in German even if some input data is in other languages." ), ChatMessage.from_user("Tell me about {{location}}"), ] response = pipe.run( data={ "prompt_builder": { "template_variables": {"location": "Berlin"}, "template": messages, } } ) print(response["llm"]["replies"][0]) print(response["tracer"]["trace_url"]) print(response["tracer"]["trace_id"])示例要点:
- 环境变量设置必须出现在
from haystack...导入语句之前,理由如前所述; LangfuseConnector("Chat example")的第一个位置参数name是必填项,它作为该次 trace 在 Langfuse 仪表盘上的标识名;pipe.run()的返回结果中,response["tracer"]携带该组件自己的输出,包含trace_url与trace_id。
进阶:在 Agent 管道中追踪工具调用过程
Agent(智能体)场景下,管道中会包含多轮"思考 → 调用工具 → 观察结果"的循环,追踪的价值更大。下面的示例来自 LangfuseConnector 组件文档,它把LangfuseConnector与带工具(天气查询、四则运算)的Agent组合,并通过invocation_context为本次调用附加自定义上下文标记:
import os os.environ["LANGFUSE_HOST"] = "https://cloud.langfuse.com" os.environ["HAYSTACK_CONTENT_TRACING_ENABLED"] = "true" from typing import Annotated from haystack.components.agents import Agent from haystack.components.generators.chat import OpenAIChatGenerator from haystack.dataclasses import ChatMessage from haystack.tools import tool from haystack import Pipeline from haystack_integrations.components.connectors.langfuse import LangfuseConnector @tool def get_weather(city: Annotated[str, "The city to get weather for"]) -> str: """Get current weather information for a city.""" weather_data = { "Berlin": "18°C, partly cloudy", "New York": "22°C, sunny", "Tokyo": "25°C, clear skies", } return weather_data.get(city, f"Weather information for {city} not available") @tool def calculate( operation: Annotated[str, "Mathematical operation: add, subtract, multiply, divide"], a: Annotated[float, "First number"], b: Annotated[float, "Second number"], ) -> str: """Perform basic mathematical calculations.""" if operation == "add": result = a + b elif operation == "subtract": result = a - b elif operation == "multiply": result = a * b elif operation == "divide": if b == 0: return "Error: Division by zero" else: result = a / b else: return f"Error: Unknown operation '{operation}'" return f"The result of {a} {operation} {b} is {result}" if __name__ == "__main__": chat_generator = OpenAIChatGenerator() agent = Agent( chat_generator=chat_generator, tools=[get_weather, calculate], system_prompt="You are a helpful assistant with access to weather and calculator tools. Use them when needed.", exit_conditions=["text"], ) langfuse_connector = LangfuseConnector("Agent Example") pipe = Pipeline() pipe.add_component("tracer", langfuse_connector) pipe.add_component("agent", agent) response = pipe.run( data={ "agent": { "messages": [ ChatMessage.from_user("What's the weather in Berlin and calculate 15 + 27?"), ], }, "tracer": {"invocation_context": {"test": "agent_with_tools"}}, }, ) print(response["agent"]["last_message"].text) print(response["tracer"]["trace_url"])这段代码展示了两个进阶点:
- Agent 全链路追踪:Agent 内部多轮 LLM 调用、工具选择与执行过程(
haystack.agent.step.llm等 span)都会作为子 span 落入同一条 Langfuse trace; invocation_context参数:通过pipe.run()传入"tracer": {"invocation_context": {...}},可以把外部执行框架的 run id、用户 id 等键值对附加到本次 trace 上,方便在 Langfuse 中按自己的业务维度检索。
flush 机制:数据何时真正到达 Langfuse
LangfuseConnector默认在每个组件执行完毕后立即刷新(flush)数据,且刷新会阻塞当前线程直到数据成功发送到 Langfuse。这种"组件级强制 flush"保证了数据的实时性与可靠性,代价是带来额外的同步开销。
如果你希望减少同步开销,可以把环境变量HAYSTACK_LANGFUSE_ENFORCE_FLUSH设为"false"来禁用组件级 flush。但要格外小心:禁用后如果程序崩溃,尚未刷出的追踪数据可能丢失。因此你必须保证在程序退出前显式调用langfuse.flush()。
以下是在普通脚本中安全退出的写法:
from haystack.tracing import tracer try: # your code here finally: tracer.actual_tracer.flush()也可以借助 FastAPI 的关闭事件处理器:
from haystack.tracing import tracer # ... @app.on_event("shutdown") async def shutdown_event(): tracer.actual_tracer.flush()注意上面两段代码都通过tracer.actual_tracer.flush()调用——这是因为haystack.tracing.tracer是一个ProxyTracer容器(见 haystack/tracing/tracer.py),真正的 Langfuse 追踪器存放在actual_tracer属性中,flush方法由它提供。
深入 API:LangfuseConnector 的构造参数与运行契约
LangfuseConnector.__init__的完整签名如下:
__init__( name: str, public: bool = False, public_key: Secret | None = Secret.from_env_var("LANGFUSE_PUBLIC_KEY"), secret_key: Secret | None = Secret.from_env_var("LANGFUSE_SECRET_KEY"), httpx_client: httpx.Client | None = None, span_handler: SpanHandler | None = None, *, host: str | None = None, langfuse_client_kwargs: dict[str, Any] | None = None ) -> None各参数语义如下:
| 参数 | 类型 | 说明 |
|---|---|---|
name | str(必填) | trace 的名称,用于在 Langfuse 仪表盘上标识该次追踪运行 |
public | bool,默认False | 追踪数据是否公开。为True时,任何持有 trace URL 的人都能访问;为False时仅 Langfuse 账号所有者可见 |
public_key | Secret \| None | Langfuse 公钥,默认从LANGFUSE_PUBLIC_KEY环境变量读取 |
secret_key | Secret \| None | Langfuse 私钥,默认从LANGFUSE_SECRET_KEY环境变量读取 |
httpx_client | httpx.Client \| None | 用于 Langfuse API 调用的自定义 HTTPX 客户端。注意:从 YAML 反序列化管道时,自定义客户端会被丢弃,由 Langfuse 自建默认客户端(HTTPX 客户端无法序列化) |
span_handler | SpanHandler \| None | 自定义 span 处理器;不传则使用DefaultSpanHandler。它控制 span 的创建方式与创建后的处理逻辑 |
host | str \| None(仅关键字参数) | Langfuse API 地址,也可通过LANGFUSE_HOST环境变量设置,默认https://cloud.langfuse.com |
langfuse_client_kwargs | dict[str, Any] \| None(仅关键字参数) | 传给 Langfuse 客户端的额外配置项字典,可用于自定义客户端行为 |
run方法的签名与返回契约:
run(invocation_context: dict[str, Any] | None = None) -> dict[str, str]- 参数
invocation_context:本次调用的附加上下文字典。适合把外部执行框架的 run id、用户 id 等信息标记到本次调用上,这些键值对会出现在 Langfuse trace 中; - 返回值:包含三个键的字典——
name:追踪组件的名称;trace_url:追踪数据的访问 URL;trace_id:trace 的 ID。
此外,LangfuseConnector实现了to_dict()/from_dict()用于组件序列化与反序列化,这使它能够被纳入 Haystack 的 YAML 管道定义体系,支持管道整体导出与复现(自定义httpx_client在反序列化时会被丢弃,但span_handler会通过其自身的to_dict/from_dict参与序列化)。
高级定制:用 SpanHandler 深度控制 span 的创建与加工
SpanHandler是langfuse-haystack集成中最重要的扩展点(抽象基类,见关联文档haystack_integrations.tracing.langfuse.tracer一节)。它定义了两个关键扩展方法:
create_span(context: SpanContext) -> LangfuseSpan:决定创建什么类型的 span。基于SpanContext的默认逻辑是:- 如果没有父 span,则创建一条新 trace;
- 对 LLM 类组件创建generation span(专门承载模型调用信息);
- 对其他组件创建default span。
handle(span: LangfuseSpan, component_type: str | None) -> None:在组件执行完毕、span 被产出后处理该 span。官方列出的典型用途包括:抽取并附加 Token 用量统计、补充模型信息、记录时间指标(如首 Token 延迟)、设置日志级别用于质量监控、添加自定义指标与观测数据。
SpanContext封装了创建 span 所需的全部上下文信息,其字段如下:
| 字段 | 说明 |
|---|---|
name | 要创建的 span 名称,对组件而言通常是组件名 |
operation_name | 被追踪的操作名(如haystack.pipeline.run),用于判断是否应无警告地新建 trace |
component_type | 创建 span 的组件类型(如OpenAIChatGenerator),用于决定 span 类型 |
tags | 附加到 span 的元数据,包含组件输入/输出数据等追踪信息 |
parent_span | 父 span;为None时创建新 trace |
trace_name | 创建父 span 时使用的 trace 名称,默认为"Haystack" |
public | trace 是否公开可见,默认False |
实现自定义处理器最常用的方式是继承DefaultSpanHandler并覆写handle。下面是关联文档中的基础示例:
from haystack_integrations.tracing.langfuse import DefaultSpanHandler, LangfuseSpan from typing import Optional class CustomSpanHandler(DefaultSpanHandler): def handle(self, span: LangfuseSpan, component_type: Optional[str]) -> None: # Custom span handling logic, customize Langfuse spans however it fits you # see DefaultSpanHandler for how we create and process spans by default pass connector = LangfuseConnector(span_handler=CustomSpanHandler())再结合 LangfuseConnector 组件文档 中的示例,你可以访问 span 内部的追踪数据与底层 span 对象,实现"质量告警"这类业务逻辑——例如检测 LLM 回复过短并给该 span 打上WARNING级别:
from haystack_integrations.tracing.langfuse import ( LangfuseConnector, DefaultSpanHandler, LangfuseSpan, ) from typing import Optional class CustomSpanHandler(DefaultSpanHandler): def handle(self, span: LangfuseSpan, component_type: Optional[str]) -> None: # Custom logic to add metadata or modify span if component_type == "OpenAIChatGenerator": output = span._data.get("haystack.component.output", {}) if len(output.get("text", "")) < 10: span._span.update(level="WARNING", status_message="Response too short") ## Add the custom handler to the LangfuseConnector connector = LangfuseConnector(span_handler=CustomSpanHandler())这里的span._data保存着该 span 关联的追踪数据(key 与 Haystack 管道打点时的 tag 一一对应,如haystack.component.output),而span._span是 Langfuse 客户端侧的底层 span 对象,可以直接调用 Langfuse 的update()方法修改级别与状态信息。SpanHandler还提供init_tracer(tracer)钩子,由LangfuseTracer在内部调用,把 Langfuse 客户端实例注入处理器。
内部机制:LangfuseTracer 与 LangfuseSpan 如何桥接两个生态
LangfuseTracer是 HaystackTracer接口的 Langfuse 实现,充当两个生态的桥接层,构造参数如下:
__init__( tracer: langfuse.Langfuse, name: str = "Haystack", public: bool = False, span_handler: SpanHandler | None = None, ) -> Nonetracer:Langfuse 客户端实例;name:管道或组件名称,用于在 Langfuse 仪表盘上标识追踪运行,默认"Haystack";public:追踪数据是否公开,默认False;span_handler:自定义 span 处理器,默认使用DefaultSpanHandler。
它对外提供五个关键方法,对应 haystack/tracing/tracer.py 中Tracer抽象基类的能力:
trace(operation_name, tags, parent_span):作为上下文管理器创建并管理一个 span(实现自 HaystackTracer.trace契约);flush():把所有挂起的 span 一次性刷新到 Langfuse;current_span():返回当前活跃 span,无则返回None;get_trace_url():返回追踪数据的 URL;get_trace_id():返回 trace ID。
与之配套的LangfuseSpan则是 HaystackSpan接口的 Langfuse 实现(基类定义见 haystack/tracing/tracer.py#L14-L74),构造时接收由langfuse.get_client().start_as_current_observation创建的上下文管理器。它实现的核心方法包括:
| 方法 | 行为 |
|---|---|
set_tag(key, value) | 设置普通标签 |
set_content_tag(key, value) | 设置内容类标签。内容指查询、文档、回答等敏感信息,Haystack 默认关闭该行为,需通过HAYSTACK_CONTENT_TRACING_ENABLED=true开启(参见 haystack/tracing/tracer.py#L49-L66) |
raw_span() | 返回底层的 Langfuse span 实例(LangfuseClientSpan) |
get_data() | 返回与该 span 关联的数据字典 |
get_correlation_data_for_logs() | 返回用于日志关联增强的相关性数据 |
值得补充的是,所有写入 span 的 tag 值都会经过类型收敛(coercion)处理:coerce_tag_value(见 haystack/tracing/utils.py)会把非基础类型的值序列化为 JSON 字符串,确保追踪后端兼容。
下图展示了在 Langfuse 追踪详情页中,一次管道运行里 LLM generation span 的典型形态(输入提示词、输出、元数据、延迟与成本信息一目了然),图片来自 docs-website/versioned_docs/version-2.23/development/tracing.mdx:
序列化与迁移:把 Langfuse 集成纳入管道工程化体系
LangfuseConnector与SpanHandler都实现了to_dict()/from_dict(),意味着追踪配置可以完整地进入 Haystack 的 YAML 管道序列化体系,实现"管道定义即代码、可版本化、可复现"。两条实践要点:
- 自定义
httpx_client无法序列化,管道从 YAML 反序列化时该配置会被丢弃,Langfuse 会自行创建默认客户端——因此生产环境建议依赖LANGFUSE_*环境变量而非代码内客户端配置; - 自定义
span_handler通过自身的to_dict/from_dict参与序列化,这让质量告警、指标增强等定制逻辑可以随管道配置一起被保存与复现。
在仓库的 MIGRATION.md 中,还保留着完整的迁移示例(需要先执行pip install langfuse-haystack),展示了LangfuseConnector与DefaultSpanHandler、SpanContext、ObservationSpanType等追踪模块的配合用法,可作为从旧版本迁移或深度定制的参考起点。
小结
LangfuseConnector以"零连线伴跑"的设计,把 Haystack 管道运行的全链路可观测性低成本地接入 Langfuse:一条环境变量开启内容追踪,一个组件实例完成 trace 上报,run()返回的trace_url/trace_id即可直达 Langfuse 详情页。在需要深度定制的场景,SpanHandler抽象提供了 span 创建(trace / generation / default)与加工(Token 统计、耗时记录、质量告警、自定义元数据)两个层次的扩展点。无论是常规 RAG 管道、多工具 Agent,还是需要工程化序列化的生产部署,这套组合都能提供从"能看见"到"可定制、可复现"的完整追踪能力。
【免费下载链接】haystackOpen-source AI orchestration framework for building context-engineered, production-ready LLM applications. Design modular pipelines and agent workflows with explicit control over retrieval, routing, memory, and generation. Built for scalable agents, RAG, multimodal applications, semantic search, and conversational systems.项目地址: https://gitcode.com/GitHub_Trending/ha/haystack
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考