
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/haystackLangfuse 是 LLM 应用的可观测性平台而 Haystack 通过LangfuseConnector这个零连线组件把管道内每一次运行的完整链路——提示词、模型调用、组件输入输出、Token 消耗与耗时——自动送入 Langfuse 仪表盘。读完本文你将掌握 LangfuseConnector 的安装配置、环境变量语义、管道与 Agent 两种接入方式、flush 数据刷新策略以及通过SpanHandler深度定制追踪内容的高级玩法。为什么要在 Haystack 管道中接入 LangfuseLangfuseConnector的作用是把 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.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.comHAYSTACK_LANGFUSE_ENFORCE_FLUSH可选设为false可关闭每个组件执行后的强制 flush详见下文 flush 一节重要必须先设置环境变量再导入任何 Haystack 组件。这是因为 Haystack 在导入过程中就会初始化内部的追踪组件参见 haystack/tracing/tracer.py#L118 中is_content_tracing_enabled的读取逻辑。更稳妥的做法是在运行脚本前于 shell 中设置这些变量把配置与代码分离便于不同环境的管理。快速上手在 RAG 对话管道中接入 Langfuse 追踪下面是一个最小可用示例用ChatPromptBuilderOpenAIChatGenerator构建一条对话管道并把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(modelgpt-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, fWeather 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 fError: Unknown operation {operation} return fThe result of {a} {operation} {b} is {result} if __name__ __main__: chat_generator OpenAIChatGenerator() agent Agent( chat_generatorchat_generator, tools[get_weather, calculate], system_promptYou 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(Whats 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 traceinvocation_context参数通过pipe.run()传入tracer: {invocation_context: {...}}可以把外部执行框架的 run id、用户 id 等键值对附加到本次 trace 上方便在 Langfuse 中按自己的业务维度检索。flush 机制数据何时真正到达 LangfuseLangfuseConnector默认在每个组件执行完毕后立即刷新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方法由它提供。深入 APILangfuseConnector 的构造参数与运行契约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.comlangfuse_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追踪数据的访问 URLtrace_idtrace 的 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用于判断是否应无警告地新建 tracecomponent_type创建 span 的组件类型如OpenAIChatGenerator用于决定 span 类型tags附加到 span 的元数据包含组件输入/输出数据等追踪信息parent_span父 span为None时创建新 tracetrace_name创建父 span 时使用的 trace 名称默认为Haystackpublictrace 是否公开可见默认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_handlerCustomSpanHandler())再结合 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(levelWARNING, status_messageResponse too short) ## Add the custom handler to the LangfuseConnector connector LangfuseConnector(span_handlerCustomSpanHandler())这里的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, ) - NonetracerLangfuse 客户端实例name管道或组件名称用于在 Langfuse 仪表盘上标识追踪运行默认Haystackpublic追踪数据是否公开默认Falsespan_handler自定义 span 处理器默认使用DefaultSpanHandler。它对外提供五个关键方法对应 haystack/tracing/tracer.py 中Tracer抽象基类的能力trace(operation_name, tags, parent_span)作为上下文管理器创建并管理一个 span实现自 HaystackTracer.trace契约flush()把所有挂起的 span 一次性刷新到 Langfusecurrent_span()返回当前活跃 span无则返回Noneget_trace_url()返回追踪数据的 URLget_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_ENABLEDtrue开启参见 haystack/tracing/tracer.py#L49-L66raw_span()返回底层的 Langfuse span 实例LangfuseClientSpanget_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),仅供参考