Langfuse 内部可观测性实践:基于 OpenTelemetry 的统一埋点、指标与追踪体系解析
【免费下载链接】langfuse🪢 Open source AI engineering platform: LLM evals, observability, metrics, prompt management, playground, datasets. Integrates with OpenTelemetry, LangChain, OpenAI SDK, LiteLLM, and more. 🍊YC W23项目地址: https://gitcode.com/GitHub_Trending/la/langfuse
导读
本文围绕 packages/shared/src/server/instrumentation/README.md 展开,系统讲解 Langfuse 自研代码库中如何**最大化使用 OpenTelemetry(Otel)**进行应用自身可观测性建设:从instrumentAsync/instrumentSync埋点函数、recordGauge/recordIncrement等指标记录 API,到 web 与 worker 两个进程的 SDK 初始化配置,再到 trpc 中间件与自动埋点清单。读完本文,你将掌握 Langfuse 内部"以 Otel 为统一抽象、可自由切换可观测后端"的设计思路,并能直接在业务代码中复用它提供的埋点与指标原语。
一、核心设计动机:为什么"尽可能使用 Otel"
文档开篇即给出了 Langfuse 内部可观测性的总原则:
Throughout our applications we want to use as much Otel as possible.
这一原则的背后有三个现实收益(文档原文 + 源码印证):
- 可观测后端可灵活更换:统一使用 Otel 标准后,trace、metrics 可以被导出到任意兼容后端,不会被单一厂商锁定;
- 享受 Otel 社区生态:自动埋点、SDK、Exporter、资源探测器等都由社区维护,团队无需自研协议层;
- 标准化的上下文传播:span 与 baggage 的传播遵循 Otel 规范,便于在多服务、多进程间串联。
从实现看,这套抽象落在 packages/shared/src/server/instrumentation/index.ts 中,它是@langfuse/shared对外暴露的可观测性入口,同时被 web 与 worker 两个应用消费。文档还明确了一条工程约定:"When building adding new infrastructure, we should search for auto instrumentations for our code base"——每接入一块新基础设施(如新的 Redis 客户端、新的 HTTP 客户端),优先查找对应的官方自动埋点包,而不是手写埋点。
二、埋点 API:用instrumentAsync/instrumentSync包裹函数
文档给出的第一个使用方式是:
Use the
instrumentorinstrumentfunctions to wrap your functions with Otel instrumentation. This will automatically create spans for your functions and send them to the Otel collector. If an instrumented function throws, exceptions will be added to the span and the span will be marked as failed.
需要说明的是,仓库实际导出的函数名为instrumentAsync与instrumentSync(README 中简写为 instrument)。二者分别用于异步与同步回调,签名一致:
// 摘自 packages/shared/src/server/instrumentation/index.ts export async function instrumentAsync<T>( ctx: SpanCtx, callback: AsyncCallbackFn<T>, // (span: opentelemetry.Span) => Promise<T> ): Promise<T>; export function instrumentSync<T>( ctx: SpanCtx, callback: SyncCallbackFn<T>, // (span: opentelemetry.Span) => T ): T;两个函数都会在回调执行期间创建并激活一个全新的 Otel span,并把回调结果原样返回给调用方。使用时你可以在回调内拿到当前 span 对象,直接对它设置属性。
SpanCtx 参数详解
| 参数 | 类型 | 作用 |
|---|---|---|
name | string | span 名称,会作为该段操作在 trace 视图中的标识 |
spanKind | opentelemetry.SpanKind | span 类型(如CLIENT、SERVER、PRODUCER、CONSUMER、INTERNAL),按 Otel SpanKind 规范 取值 |
rootSpan | boolean | 是否作为 root span(即切断父 trace 关系),参见 Otel trace 语义约定 |
traceScope | string | 指定使用的 tracer 名称,未指定时回退到callback.name(函数名) |
traceContext | TCarrier({ traceparent?, tracestate? }) | 通过 W3C Trace Context 载体恢复外部传入的 trace 上下文,用于跨服务续接链路 |
startNewTrace | boolean | 是否开启全新 trace,切断一切父 span 关系(但会保留 baggage) |
源码级实现原理
以 index.ts 中instrumentAsync为例,其执行流程为:
- 确定激活上下文:按优先级取
startNewTrace→ 在ROOT_CONTEXT上重建 baggage(切断父 trace 但保留 baggage,用于跨新 trace 传递业务标签);否则若传入traceContext→ 用opentelemetry.propagation.extract解析出上下文;否则直接用当前活跃上下文; - 开启 span:
getTracer(ctx.traceScope ?? callback.name).startActiveSpan(ctx.name, { root, kind }, activeContext, ...)。root的取值是startNewTrace || (!traceContext && rootSpan),即显式开新 trace 或声明为 root 且无外部上下文时都会成为根 span; - 搬运 baggage:把当前 baggage 中的全部键值对写入 span attribute(
span.setAttribute(k, v.value)),这样下游查询(如 ClickHouse 查询标签)可以沿 baggage 传播的上下文工作; - 异常处理:回调抛错时调用
traceException(ex, span)记录异常、打上错误标签、将 span 状态置为ERROR,然后重新抛出异常,不吞错、不改变调用方语义。
配套的 span 辅助函数
同一文件还导出了一组便捷工具:
getCurrentSpan():返回当前活跃 span(opentelemetry.trace.getActiveSpan()),用于在任意深度的代码中安全地获取当前 span;addTagsToCurrentSpan(attributes):向当前活跃 span 批量追加属性;traceException(ex, span?, code?):把异常记录为 Otel span event(recordException),同时为 Datadog 错误追踪设置error.stack、error.message、error.type三个标签,并将 span 状态置为ERROR。它对异常做了兼容处理:Error实例取message/name/stack,普通对象则JSON.stringify兜底;addUserToSpan(attributes, span?):把userId、projectId、orgId、plan、apiKeyId、publicKey等业务属性同时写入span attribute 与 baggage,例如user.id、langfuse.project.id、langfuse.org.plan、langfuse.api_key.id。测试 index.test.ts 专门验证了"不会把用户 email 写入 span 或 baggage"这一隐私约束。
三、指标记录 API:Gauge / Counter / Histogram
文档给出的第二类原语是:
Use
recordGauge,recordCounter,recordHistogramto record metrics. These will be sent to the Otel collector.
仓库实际导出的函数名为recordGauge、recordIncrement(即 Counter 语义)、recordHistogram,另有补充的recordDistribution。它们的统一签名是(stat: string, value?: number, tags?: { [tag: string]: string | number }),并采用双通道输出:
- Datadog:无条件调用
dd.dogstatsd.gauge/increment/histogram/distribution,指标进入 Datadog; - AWS CloudWatch(可选):仅当环境变量
ENABLE_AWS_CLOUDWATCH_METRIC_PUBLISHING === "true"时启用(见 env.ts 附近的配置定义)。
CloudWatch 通道的实现细节非常值得借鉴:
- 30 秒批量冲刷:指标先写入进程内缓存
metricCache,每 30 秒由flushMetricsToCloudWatch一次性通过PutMetricDataCommand发送(index.ts),避免高频打点造成过多 API 调用;失败仅logger.warn告警,不阻塞业务; - 命名空间固定为
Langfuse; - 标签扁平化:以
.depth、.rate、.dlq_oldest_age结尾的指标,其 tags 会被按 key 字典序排序后拼进 CloudWatch 的指标名(排除unit),例如langfuse.queue.depth会展开成langfuse.queue.depth.surface_worker.route_langfuse.queue.monitor这样的复合名称,其余指标不受影响。
此外还提供convertQueueNameToMetricName(queueName)把队列名转换为 Datadog 指标名:legacy-ingestion-queue→langfuse.queue.legacy_ingestion(规则是连字符转下划线、去掉_queue后缀、加langfuse.queue.前缀)。worker 侧消费队列时可以方便地复用该命名。
实际调用示例
- s3SlowdownTracking.ts 中使用
recordIncrement("langfuse.s3_slowdown.marked", 1)记录 S3 慢请求标记次数; - queryOutcome.ts、processEventBatch.ts、blockEvaluators.ts 等模块也广泛使用了这些指标 API,作为请求成功/失败率、耗时分布等业务指标的载体。
四、进程级配置:web 与 worker 的 instrumentation 文件
文档指出:
webandworkerhave an instrumentation.ts file, which configures otel for the application.
两个应用各有一份初始化文件,结构高度一致,都是dd.init()+NodeSDK组合:
- web:web/src/observability.config.ts(由 web/src/instrumentation.ts 的 Next.js
register()钩子按运行时与开关动态加载) - worker:worker/src/instrumentation.ts
SDK 核心配置
两份配置共享以下关键项:
| 配置项 | 取值 | 说明 |
|---|---|---|
resource | service.name=env.OTEL_SERVICE_NAME,service.version=env.BUILD_ID | 标识服务身份,便于按服务名与版本检索 |
traceExporter | OTLPTraceExporter,URL 为${env.OTEL_EXPORTER_OTLP_ENDPOINT}/v1/traces | 通过 OTLP/proto 协议把 trace 发送到 Otel Collector |
resourceDetectors | envDetector、awsEcsDetector、containerDetector | 自动从环境变量、AWS ECS 元数据、容器元数据中补充资源属性 |
sampler(仅 web) | TraceIdRatioBasedSampler(env.OTEL_TRACE_SAMPLING_RATIO) | 按比例采样控制 trace 数据量 |
自动埋点清单(instrumentations)
两份配置共同启用的自动埋点包括:
IORedisInstrumentation(携带ioredisRequestHook,见下文脱敏策略)HttpInstrumentation:要求出站 span 必须具有父 span(requireParentforOutgoingSpans: true),忽略健康检查路径(/api/public/health、/api/public/ready、/api/health)与回环地址127.0.0.1的出站请求;requestHook会把 span 名称规范化为METHOD /path形式,并将/_next/static/*聚合为通配路径、去掉/index后缀PrismaInstrumentation:忽略prisma:client:serialize、prisma:engine:query等低价值 span 类型,降低数据库层噪音AwsInstrumentation:自动埋点 AWS SDK 调用WinstonInstrumentation(disableLogSending: true):只做日志与 trace 的关联,不把日志内容发送到远端BullMQInstrumentation(useProducerSpanAsConsumerParent: true):让消费者 span 挂在生产者 span 下,串联消息队列的完整链路
worker 进程额外启用了ExpressInstrumentation与UndiciInstrumentation(对非127.0.0.1/localhost的外部请求打点并规范 span 名称)。
Redis 命令脱敏:ioredisRequestHook
bootstrap/ioredisRequestHook.ts 是防止敏感信息进入 trace 的典型实现:
AUTH、HELLO命令直接记录为[REDACTED];- 若命令首参数命中 API Key 缓存前缀
API_KEY_CACHE_KEY_PREFIX,则其余参数全部替换为[REDACTED]; - 其余命令以
CMD arg1 arg2 ...形式写入redis.full_command属性。
该 hook 从@langfuse/shared/instrumentation/bootstrap子路径导出,原因是 bootstrap 目录专门放置必须在sdk.start()之前加载、且不能引用 server barrel 的初始化代码(见 bootstrap/index.ts 的注释约定)。
调用方 SDK 识别:sdkName / sdkVersion
web 的HttpInstrumentation.requestHook会从入站请求头中提取x-langfuse-sdk-name、x-langfuse-sdk-version(sdkName.ts),并写入sdk_name、sdk_version两个 span 属性。这套机制让 Langfuse 能在 trace 中直接区分请求来自 Python SDK、JavaScript SDK 还是其他客户端,且 SDK 名称被规范化到封闭集合(python/javascript),避免攻击者通过任意字符串污染标签。
web 侧的两个进阶技巧
- Next.js 包装 span 的静音:observability.config.ts 实现了
PassthroughTracer与ScopeFilteringTracerProvider,把next.js这个 tracer scope 静音——这些 span 不再被创建,子 span 直接挂在http.server下,既减少了数据量,又保留 Next.js 在父 span 上打的路由模板http.route; - 懒加载与优雅退出:web/src/instrumentation.ts 仅在 Node.js 运行时且
NEXT_PUBLIC_LANGFUSE_RUN_NEXT_INIT未显式设为"false"时执行 init;同时安装installProcessErrorHandlers,发生致命错误时先排空在途请求(drainAndClose)再退出,避免直接 5xx。
五、trpc 链路增强:@baselime/trpc-opentelemetry-middleware
文档特别提到:
For trpc, we use
@baselime/trpc-opentelemetry-middlewareto enrich spans with trpc inputs and outputs.
在 web/src/server/api/trpc.ts 中导入tracing中间件,用于把 trpc 过程的输入(inputs)与输出(outputs)追加到对应 span 上。这样在排查"某个 trpc 调用为何失败/变慢"时,可以直接在 trace 里看到请求入参与返回结果,无需再翻应用日志。
六、行为验证:单元测试如何约束埋点语义
instrumentation/index.test.ts 用 vitest 覆盖了四个关键行为,可作为理解语义的活文档:
startNewTrace保留业务上下文:在contextWithLangfuseProps注入 worker surface/route 后,无论instrumentAsync还是instrumentSync开新 trace,ClickHouse 查询标签中的surface: "worker"、route: "langfuse.queue.monitor"、projectId都完好保留;- 公共 API 调用方属性可传播:入站请求头中的 SDK 名称/版本与 User-Agent 会进入 ClickHouse 查询标签;
- 拒绝来自入站 baggage 的伪造属性:攻击者通过 baggage 注入的
sdkVersion/userAgent不会被采纳,防止标签污染; - 隐私约束:
addUserToSpan不会把邮箱写入 span 或 baggage。
七、实践建议:接入新基础设施时的埋点 checklist
结合文档末尾约定与仓库实现,在 Langfuse 生态内新增代码时可遵循以下 checklist:
- 能用自动埋点解决的不手写:先搜索
@opentelemetry/instrumentation-*生态是否有对应客户端(HTTP、ioredis、Prisma、AWS SDK、BullMQ、Express、Undici、Winston 都是现成案例); - 业务函数级埋点用
instrumentAsync/instrumentSync,并善用traceScope归类 tracer、addTagsToCurrentSpan补属性; - 需要跨服务续接链路时,通过
traceContext(traceparent/tracestate)提取;需要独立链路时用startNewTrace(注意它保留 baggage); - 统计类指标用
recordIncrement/recordGauge/recordHistogram,命名遵循langfuse.前缀与convertQueueNameToMetricName约定; - 涉及凭据、API Key、用户隐私的数据绝不落入 span:参考
ioredisRequestHook的脱敏模式与addUserToSpan的隐私白名单; - 新进程接入时参考 web/src/observability.config.ts 与 worker/src/instrumentation.ts 的
NodeSDK配置模板,并注意 bootstrap 目录的加载时机约束。
结语
Langfuse 内部可观测性体系的精髓在于"标准先行、封装收敛":对外全部采用 Otel 标准(trace 与 metrics),对内通过instrumentAsync、recordGauge等少数几个函数收敛全部埋点语义,再以 web/worker 两份instrumentation配置集中管理 Exporter、采样率、自动埋点与脱敏规则。这套模式既保证了后端可替换的灵活性,也让业务代码的埋点成本降到最低——这也正是该 README 文档希望传达并持续维护的工程准则。
【免费下载链接】langfuse🪢 Open source AI engineering platform: LLM evals, observability, metrics, prompt management, playground, datasets. Integrates with OpenTelemetry, LangChain, OpenAI SDK, LiteLLM, and more. 🍊YC W23项目地址: https://gitcode.com/GitHub_Trending/la/langfuse
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考