callbacks 回调系统
不改任何业务代码,就能给每个节点埋点——日志、监控、计时、链路追踪。今天讲 5 个触发时机、Handler 接口,以及流式回调那个必须知道的坑。
回调解决什么问题
log.Print——10 个地方改 30 行,还有些是第三方组件根本改不动。上线后想临时关掉监控?又得改回去。这类"和业务无关、却要撒进每个角落"的需求,硬塞进业务代码里就是灾难。你想知道"每个节点耗时多少、输入输出是啥、有没有报错",用来做监控和排查。但你不想在每个组件里手写 log.Print——那样代码脏、还改不动第三方组件。回调让你从外部统一挂上这些逻辑。
@Aspect、Python 的装饰器是同一种思想。5 个触发时机
Eino 在每个节点执行的关键点触发回调(callbacks/、internal/callbacks):
开始前→ 节点执行中→ OnEnd
成功后 或 OnError
出错时
流式输入开始 OnEndWithStreamOutput
流式输出结束
OnStartWithStreamInput / OnEndWithStreamOutput)针对流式场景——输入/输出是流的时候用这两个。一共 5 个时机。每个节点执行都会走"开始→(成功|出错)结束",回调就挂在这些点上。这 5 个时机在源码里就是 callbacks/interface.go:117 起的 TimingOnStart / TimingOnEnd / TimingOnError / TimingOnStartWithStreamInput / TimingOnEndWithStreamOutput 枚举。Handler 接口(真实代码)
一个回调 handler 实现这些方法(callbacks/interface.go,用 Builder 构造更方便):
handler := callbacks.NewHandlerBuilder().
OnStartFn(func(ctx, info *RunInfo, input CallbackInput) context.Context {
// 节点开始:记下开始时间、打日志
return context.WithValue(ctx, startKey, now)
}).
OnEndFn(func(ctx, info *RunInfo, output CallbackOutput) context.Context {
// 节点成功:算耗时、上报监控
return ctx
}).
OnErrorFn(func(ctx, info *RunInfo, err error) context.Context {
// 节点出错:记录错误
return ctx
}).Build()
如果让你自己做,最朴素的写法大概是往每个组件里手插日志:
// 你的极简版:埋点和业务混在一起,每个组件都要改
func (m *MyModel) Generate(ctx, msgs) (*Message, error) {
log.Print("start", time.Now()) // ← 手插埋点
out, err := 真正调模型(msgs)
log.Print("end", time.Now()) // ← 又手插一遍
return out, err
}
Eino 真实版把埋点抽到框架执行层,你只在外部写一次 Handler:
handler := callbacks.NewHandlerBuilder().
OnStartFn(...).OnEndFn(...).Build() // ← 埋点集中在这,业务代码零改动
真实版多出的这层解决了什么? 极简版里每个组件都得手插、第三方组件根本插不进、想临时关掉还要改回去;真实版把埋点和业务彻底分开——写一次、挂一次、对所有节点(含第三方)生效、想摘随时摘。这就是"横切"的价值。OnStartFn 里把 time.Now() 用 context.WithValue 塞进 ctx;在 OnEndFn 里取出来算 time.Since(start) 打印。其余三个时机(OnError / 两个流式)完全不用管——这正是 NewHandlerBuilder(callbacks/handler_builder.go:109)的价值:直接实现 Handler 接口(callbacks/interface.go:85 的 Handler = callbacks.Handler)要写全 5 个方法,用 Builder 只挑你关心的时机设,剩下的框架给空实现。这个 Handler 挂上后,图里每个节点自动被计时,业务代码一行没动。怎么挂上去
两种挂法:全局(对所有执行生效)或运行时(这一次执行生效):
// 方式一:全局挂(进程启动时,对之后所有调用生效)
callbacks.AppendGlobalHandlers(handler)
// 方式二:运行时挂(只对这一次 Invoke/Stream 生效)
r.Invoke(ctx, input, compose.WithCallbacks(handler))
RunInfo:告诉你"现在在哪"
每次回调都带一个 RunInfo,告诉你当前是哪个节点、什么组件类型、叫什么名字:
type RunInfo struct {
Name string // 节点名
Type string // 具体类型
Component components.Component // 组件类别(ChatModel/Tool/...)
}
流式回调的坑(必须知道)
流式场景(OnEndWithStreamOutput)里,回调拿到的是一个 StreamReader(Day 03)。而流只能被消费一次——如果你的 handler 把流读完了,真正的业务下游就读不到了!
StreamReader.Copy)。框架/你在把流交给回调前会 Copy 一份,让回调读副本、业务读原本,互不影响。写流式 handler 时切记:要么用框架已 Copy 好的副本,要么自己 Copy——绝不能直接把业务的流读干。这是流式回调最常见的 bug 来源。另外流式 handler 读流是异步的,别在里面做重活阻塞主流程。👶 小白:我在 OnEndWithStreamOutput 里就是想打个日志,把模型输出的流读出来看看内容——这能有什么问题?
👨🏫 老师:问题大了。接着"一次性水管"的比方——模型输出这根水管里的水(chunk)流过一次就没了。你在回调里把它读完打日志,等于把水全喝光,真正等着这股水的业务下游(前端、下一个节点)就一滴都接不到,界面上表现为"模型好像没输出"。
👨🏫 解法就是这节说的 Copy:把水管分叉成两根,回调读副本、业务读原本,各喝各的互不影响。实践上通常框架已经替你 Copy 好、把副本交给回调;你自己若要再分给别处,就自己再 Copy 一份——铁律是"绝不直接读业务那根流"。另外别在回调里对着流做重活(它是异步读的),拖慢主流程就本末倒置了:回调的本分是"只观测、不打扰"。
典型用途
- 日志:每个节点的输入输出打点,排查问题。
- 指标监控:节点耗时、成功率、token 用量 → Prometheus。
- 链路追踪:给每个节点开一个 span,串成一条 trace(OpenTelemetry/Langfuse)。
- 成本核算:在模型节点的 OnEnd 里累加 token → 算这次请求花了多少钱(Day 01 成本经济学)。
今日小结 + 动手
🧠 今天你应该能回答
- 回调解决什么?为什么叫 AOP/横切?
- 5 个触发时机分别是?
- Handler 怎么写、怎么挂(全局 vs 运行时)?
- RunInfo 有什么用?
- 流式回调为什么必须 Copy 流?
✋ 动手
sed -n '1,80p' callbacks/interface.go | head -60
grep -rn 'OnStart\|OnEnd\|OnError\|WithStream' callbacks/ | head
grep -rn 'AppendGlobalHandlers\|WithCallbacks' . | head