作者:互联网 时间: 2026-08-20 08:32:56
上一篇把 Eino 的切面机制拆完了:一个 16 行的高阶函数包装,加一个挂在 ctx 里的注册表。

机制有了,这篇把它用起来——接 Langfuse,给 Agent 装一台"行车记录仪":每次 LLM 调用的输入、输出、token 数、耗时、报错,全部录下来,出事故了倒回去看。
Langfuse 是一个开源的 LLM 可观测平台(可自部署,也有云服务),核心就三层概念:
trace 一次完整请求(对应一个用户问题跑完整个图)└─ observation 图里发生的事,分两种: ├─ span 普通节点(prompt 渲染、工具调用、子图…) └─ generation LLM 调用(带模型名、token 用量——行车记录仪的主机位)Eino 对它的接入官方就有:eino-ext/callbacks/langfuse。这篇讲怎么接、接完数据长什么样、以及管道里那些不读源码不知道的坑。
全部代码就这么多(真实 API,非伪码):
package mainimport ("context""github.com/cloudwego/eino-ext/callbacks/langfuse""github.com/cloudwego/eino/callbacks")funcmain() {// 第 1 步:造 handler + flushercbh, flusher := langfuse.NewLangfuseHandler(&langfuse.Config{Host: "http://localhost:3000", // 自部署地址;云服务用 PublicKey: "pk-lf-...",SecretKey: "sk-lf-...",})// 第 2 步:挂到全局(第 87 篇的三层作用域,这里用进程级)callbacks.AppendGlobalHandlers(cbh)// 第 3 步:跑图之前给这次 trace 挂上会话身份ctx := langfuse.SetTrace(context.Background(),langfuse.WithName("demo-agent"),langfuse.WithSessionID("session-123"),langfuse.WithUserID("user-456"),)_ = ctx // → runnable.Invoke(ctx, "问题"),图的每个节点自动被录下来defer flusher() // 第 4 步(容易忘):退出前把队列里的事件冲出去}四个关键点:
SetTrace 挂的是 ctx 不是全局。sessionID/userID/tags 这些身份跟一次请求走,同一进程并发跑多个用户互不串线——和第 87 篇"注册表挂 ctx"是同一个设计判断。AppendGlobalHandlers 只能在启动时调一次(第 87 篇讲过,非线程安全)。flusher() 必须调。事件不是实时发的,是攒批异步发的(下文详述),进程退出前不 flush,最后一批就丢了。ServiceName: "eino-app",但 Config 结构体里根本没有这个字段(README_zh.md:43)。文档漂移,抄示例的话以 langfuse.go:34-126 的字段为准。跑完打开 Langfuse 界面,就能看到一棵 trace 树:每个节点的输入输出、generation 里的模型名和 token 数、报错标红。这就是行车记录仪的回放界面。
CallbackHandler 实现了第 87 篇讲的 Handler 接口,翻译规则就一张表:
| Eino 侧发生什么 | Langfuse 侧创建什么 |
|---|---|
| 第一次 OnStart(图开始) | trace-create(懒创建,SetTrace 的身份此刻生效) |
| 任意节点 OnStart | span-create 或 generation-create |
| 节点 OnEnd | 对应的 update(带 output;generation 带 token 用量) |
| 节点 OnError | update + Level=ERROR,错误文本当 output |
| 组件是 ChatModel | 用 generation 而不是 span(能带 Model、TokenUsage) |
两个值得单独说的设计:
图本身也是一个 span。 整个 Runnable 执行时,框架的包装函数先对"图"触发一次 OnStart,再对每个节点触发。所以树不是"trace → 三个平级节点",而是:
trace└─ span(graph) ← 图整体 ├─ span(prompt) ├─ generation(model) ← LLM 调用,带 model + tokens └─ span(tool)父子关系靠 state 挂 ctx。 handler 内部有个两字段结构:
type langfuseState struct {traceID string// 这条 trace 是谁observationID string// 我自己是谁(子节点拿它当 parent)}OnStart 创建 span 后,把新的 state 塞进返回的 ctx;框架把这个 ctx 传进图内部,子节点的 OnStart 取出来,ParentObservationID = state.observationID——树的边就这么连上了。和第 87 篇 RunInfo 借道 ctx 完全同构:观测身份永远跟着 ctx 走。
场景 A 的真实输出:
====== 场景 A:三节点图 → trace 树 + 分批上传 ====== 图输出: T(M(P(问题))) Langfuse 收到 3 个批次(FlushAt=3): 批 1(3 条): trace-create span-create span-create 批 2(3 条): span-update generation-create generation-update 批 3(3 条): span-create span-update span-update 还原后的 trace 树: trace trace-1 name=demo-agent input="session=s-1 user=u-42" └─ span obs-3 name=graph └─ span obs-5 name=prompt ↑ end obs-5 out="P(问题)" └─ generation obs-8 name=model model=deepseek-chat ↑ end obs-8 out="M(P(问题))" tokens=128+64 └─ span obs-11 name=tool ↑ end obs-11 out="T(M(P(问题)))" ↑ end obs-3 out="T(M(P(问题)))"9 个事件还原出完整的树:1 条 trace、图级 span、三个子 observation,generation 独享 model 名和 tokens=128+64。tok 数从哪来?第 87 篇讲的组件自管路径——openai 适配器在 OnEnd 里带上 model.CallbackOutput.TokenUsage,langfuse handler 在这里消费它。两篇的机制在这里合流。
如果每个节点执行完就同步发一个 HTTP 请求给 Langfuse,观测系统会拖慢 Agent,Langfuse 挂了 Agent 也跟着挂。所以整条管道是异步的(libs/acl/langfuse):
节点 OnStart/OnEnd │ push(非阻塞,队列默认 100 格) ▼queue(chan,满了直接丢) │ consumer goroutine 攒批 ▼批:凑够 FlushAt=15 条 或 等 FlushInterval=500ms 或 总量 2.5MB │ 逐条过四道工序:采样 → 媒体处理 → 脱敏 → 截断 ▼POST /api/public/ingestion(Basic 认证 pk:sk) └─ 指数退避重试,初始 1s,最多 3 次四道工序里三道有讲究:
采样是确定性的,且按 trace 不按事件。sha256(traceID) 前 8 位十六进制转数字,除以 0xFFFFFFFF 归一化,小于采样率就留。用 hash 而不是随机数的好处:同一条 trace 的所有事件永远得到同一个判定——要么整条 trace 都在,要么整条都不在,不会出现"看到了 LLM 调用却看不到它前面的 prompt"这种半截尸体。demo 场景 B 验证:
====== 场景 B:采样按 traceID 确定性判定 ====== 20 条 trace × 判定两次:零不一致;30% 采样率下命中 5 条截断从最大的字段开始清。 单事件超过 MaxEventSizeBytes(默认 1MB)时,把 input/output/metadata 三个字段按大小排序,从最大的开始逐个清空成占位文本,直到回到预算内——不是拦腰砍,是整字段牺牲,保证剩下的字段仍是完整合法的 JSON。demo 场景 C:
====== 场景 C:大事件截断(从最大的字段开始清)====== 截断前: input=60 字节 output=200 字节 总=260(上限 100) 截断后: input=60 字节 output="(truncated)" 总=71脱敏在上传前最后一站。MaskFunc func(string) string 在消费侧对将要出门的 input/output 做替换。demo 场景 D:
====== 场景 D:MaskFunc 上传前脱敏 ====== 生产侧原文: api_key=sk-live-9f8e7d6c Langfuse 收到: input="api_key=***" output="ok"注意它脱的是出门的数据,生产侧内存里还是原文——业务逻辑完全无感。PII 脱敏的完整方案是第 90 篇的主题,这里是它的管道挂点。
重试的分类逻辑容易读漏:指数退避重试只对网络错误、5xx、429 生效;4xx(除 429)被视为永久失败,直接吞掉(consumer.go:314 的 return nil——把错误吃掉让 backoff 停止)。事件格式错了重试一百次也不会对,这个分类是对的,但代价是:如果你发的事件被 Langfuse 拒收(比如字段超限),日志里只有一行 upload error,数据静默消失。
flush() 的语义是 q.join()——用 sync.Cond 等 unfinished 归零,即阻塞到队列里所有事件都上传完。所以它既是兜底也是同步点:长跑服务可以在每轮对话结束后调一次,不必等进程退出。Config.Timeout 不填就是 0,http.Client{Timeout: 0} = 永不超时。Langfuse 服务端若僵死,重试会以"无限等待 × 3 次"的方式卡住 consumer goroutine(不至于卡业务,但事件开始堆积直至队列满)。生产上把 Timeout 显式设上。put 是非阻塞的,满了直接返回 "event send queue is full",handler 打一行日志继续跑。更隐蔽的是连锁反应:如果被丢的是 span-create,OnEnd 找不到 state,日志里会出现一串 no state in context——看到这行先怀疑队列满,不是 bug 在别处。demo 场景 E:====== 场景 E:队列满 → 事件被丢弃(非阻塞、不 panic)====== [drop] event send queue isfull: span-create n2 ...(共 6 条 drop 日志) 塞 8 个事件进 2 格队列:丢弃 6 个,进程未阻塞未 panicOnEndWithStreamOutput 里起一个 goroutine 把流 Recv 到 EOF,拼接成完整消息再补发 generation-update,defer 里 close。这意味着流式场景下 generation 的 output 是流结束后才补录的,界面上会晚到一步,属正常。EndTrace 复用的是 trace-create 事件类型(libs/acl/langfuse/langfuse.go:127-133),不是单独的 update 类型——Langfuse 的 ingestion 按 ID upsert,所以"结束 trace"实际是"带同 ID 再 create 一次"。读日志看到两条 trace-create 不要以为创建了两次。| 问题 | 答案 | 关键源码 |
|---|---|---|
| 接入要几步 | handler + 全局挂载 + SetTrace + flusher,四步 | README_zh.md(但示例有漂移) |
| 树怎么连 | state{traceID, obsID} 挂 ctx,子节点取 parent | callbacks/langfuse/langfuse.go:191 |
| ChatModel 特殊在哪 | 用 generation 而非 span,带 model + token | OnStart 的 Component 分流 |
| 会不会拖慢 Agent | 不会:非阻塞队列 + 批量 + 异步 consumer | queue.go + consumer.go |
| 会不会丢数据 | 会:队列满丢、4xx 永久失败丢、不 flush 丢 | 三处各有日志,无一 panic |
几条设计判断:
下一篇进到桥接层的源码细节:client.go 的 batch ingestion 协议、多模态消息里的 base64 图片怎么三步上传成 Langfuse media、以及这套集成怎么用 mock 做测试。