Haystack 接入 Langfuse 全流程指南:用 LangfuseConnector 追踪 LLM 管道运行
2026/9/15 18:13:29 网站建设 项目流程

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 维护了一个全局tracerProxyTracer实例),并通过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.namehaystack.component.type、输入输出规格等标签。

LangfuseConnector本质上是在管道启动时把自己注册为全局 tracer(对应LangfuseTracer,见下文"内部机制"一节),于是上述所有 span 都会被自动转发到 Langfuse,这就是它"自动追踪所有管道操作"的原理。当管道运行完成后,run()方法会返回一个包含nametrace_urltrace_id三个键的字典,方便你把追踪链接打印出来或接入手工审计流程。

环境准备与安装

使用LangfuseConnector前需要准备以下内容:

  1. 拥有一个可用的 Langfuse 账号(Langfuse 云服务或自托管实例均可);
  2. 在账号设置中获取公钥与私钥;
  3. 安装集成包:
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_urltrace_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

各参数语义如下:

参数类型说明
namestr(必填)trace 的名称,用于在 Langfuse 仪表盘上标识该次追踪运行
publicbool,默认False追踪数据是否公开。为True时,任何持有 trace URL 的人都能访问;为False时仅 Langfuse 账号所有者可见
public_keySecret \| NoneLangfuse 公钥,默认从LANGFUSE_PUBLIC_KEY环境变量读取
secret_keySecret \| NoneLangfuse 私钥,默认从LANGFUSE_SECRET_KEY环境变量读取
httpx_clienthttpx.Client \| None用于 Langfuse API 调用的自定义 HTTPX 客户端。注意:从 YAML 反序列化管道时,自定义客户端会被丢弃,由 Langfuse 自建默认客户端(HTTPX 客户端无法序列化)
span_handlerSpanHandler \| None自定义 span 处理器;不传则使用DefaultSpanHandler。它控制 span 的创建方式与创建后的处理逻辑
hoststr \| None(仅关键字参数)Langfuse API 地址,也可通过LANGFUSE_HOST环境变量设置,默认https://cloud.langfuse.com
langfuse_client_kwargsdict[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 的创建与加工

SpanHandlerlangfuse-haystack集成中最重要的扩展点(抽象基类,见关联文档haystack_integrations.tracing.langfuse.tracer一节)。它定义了两个关键扩展方法:

  1. create_span(context: SpanContext) -> LangfuseSpan:决定创建什么类型的 span。基于SpanContext的默认逻辑是:
    • 如果没有父 span,则创建一条新 trace
    • 对 LLM 类组件创建generation span(专门承载模型调用信息);
    • 对其他组件创建default span
  2. 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"
publictrace 是否公开可见,默认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, ) -> None
  • tracer: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 集成纳入管道工程化体系

LangfuseConnectorSpanHandler都实现了to_dict()/from_dict(),意味着追踪配置可以完整地进入 Haystack 的 YAML 管道序列化体系,实现"管道定义即代码、可版本化、可复现"。两条实践要点:

  • 自定义httpx_client无法序列化,管道从 YAML 反序列化时该配置会被丢弃,Langfuse 会自行创建默认客户端——因此生产环境建议依赖LANGFUSE_*环境变量而非代码内客户端配置;
  • 自定义span_handler通过自身的to_dict/from_dict参与序列化,这让质量告警、指标增强等定制逻辑可以随管道配置一起被保存与复现。

在仓库的 MIGRATION.md 中,还保留着完整的迁移示例(需要先执行pip install langfuse-haystack),展示了LangfuseConnectorDefaultSpanHandlerSpanContextObservationSpanType等追踪模块的配合用法,可作为从旧版本迁移或深度定制的参考起点。

小结

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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询