Day 16 / 共 20 天 · 第 4 周 横切与生态

callbacks 回调系统

不改任何业务代码,就能给每个节点埋点——日志、监控、计时、链路追踪。今天讲 5 个触发时机、Handler 接口,以及流式回调那个必须知道的坑。

📍 你在整门课的位置 · 第 4 周 进阶与生态(共 4 周 · 20 天)
D16 callbacks D17 流式深水 D18 类型/泛型 D19 eino-ext D20 构建收官
L01

回调解决什么问题

🤔 痛点:想给每个组件加日志/耗时统计,又不想改一堆代码怎么办? 你的图里有 10 个节点,现在要给每个都加"进入时打日志、退出时算耗时、报错时上报告警"。老办法是钻进每个组件挨个加 log.Print——10 个地方改 30 行,还有些是第三方组件根本改不动。上线后想临时关掉监控?又得改回去。这类"和业务无关、却要撒进每个角落"的需求,硬塞进业务代码里就是灾难。

你想知道"每个节点耗时多少、输入输出是啥、有没有报错",用来做监控和排查。但你不想在每个组件里手写 log.Print——那样代码脏、还改不动第三方组件。回调让你从外部统一挂上这些逻辑。

💡 本质:横切面 / AOP——把"日志·追踪·监控"从业务里剥出来,统一从外部注入 callbacks 的本质是面向切面编程(AOP):日志、tracing、监控这些需求横跨所有节点,却不属于任何一个节点的业务逻辑。Eino 把埋点做在框架的执行层(不在组件内部),于是你只需在外部写一个 Handler、挂一次,框架执行任意节点时都会在固定时机回调你——业务代码零改动。这和 Java Spring 的 @Aspect、Python 的装饰器是同一种思想。
类比:监控摄像头 回调就像在流水线的每个工位装摄像头:不改工人的操作,就能记录"这个工位几点开始、几点结束、处理了什么、出没出错"。你想加监控,装摄像头(挂 handler)即可,不用改流水线本身。这在设计上叫 AOP(面向切面编程)——把"日志/监控"这类横切关注点从业务逻辑里剥离。
L02

5 个触发时机

Eino 在每个节点执行的关键点触发回调(callbacks/internal/callbacks):

OnStart
开始前
节点执行中 OnEnd
成功后
OnError
出错时
OnStartWithStreamInput
流式输入开始
OnEndWithStreamOutput
流式输出结束
5 个时机就像生活中快递站的进出登记本 就像生活中的快递站:每个包裹(一次节点执行)进站时记一笔"几点到、谁的、什么件"(OnStart),出站派送成功记一笔(OnEnd),要是包裹破损/丢件就记异常那本(OnError);遇到成批到的传送带包裹,就用"批量进/批量出"两本专账(两个 WithStream 时机)。登记本只是记录,一笔都不会耽误包裹本身的运送——回调也一样,只观测、不改变节点的执行结果。
前三个(OnStart / OnEnd / OnError)针对普通调用;后两个(OnStartWithStreamInput / OnEndWithStreamOutput)针对流式场景——输入/输出是流的时候用这两个。一共 5 个时机。每个节点执行都会走"开始→(成功|出错)结束",回调就挂在这些点上。这 5 个时机在源码里就是 callbacks/interface.go:117 起的 TimingOnStart / TimingOnEnd / TimingOnError / TimingOnStartWithStreamInput / TimingOnEndWithStreamOutput 枚举。
callbacks 切面(AOP)包裹住节点执行 节点执行 ChatModel / Tool …(你的业务代码) OnStart 开始前 OnStartWithStreamInput OnEnd 成功后 OnError 出错时 OnEndWithStreamOutput
5 个时机像"切面"从外部包住节点:进入前 2 个、结束后 3 个。节点自身代码毫不知情。
L03

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()
读法:每个时机一个函数,都收到 ctx + RunInfo(你在哪个节点,L05)+ 输入/输出/错误。注意每个函数返回 ctx——你可以往 ctx 里塞东西(如开始时间),下一个时机取出来用(算耗时)。用 Builder 只实现你关心的时机即可。
🔧 简化版 → 真实版:如果让你自己设计"埋点",会怎么写?

如果让你自己做,最朴素的写法大概是往每个组件里手插日志:

// 你的极简版:埋点和业务混在一起,每个组件都要改
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()   // ← 埋点集中在这,业务代码零改动
真实版多出的这层解决了什么? 极简版里每个组件都得手插、第三方组件根本插不进、想临时关掉还要改回去;真实版把埋点和业务彻底分开——写一次、挂一次、对所有节点(含第三方)生效、想摘随时摘。这就是"横切"的价值。
📝 例子:只做"节点耗时统计",只需设 OnStart / OnEnd 两个时机 想给全图做个最简单的耗时监控:在 OnStartFn 里把 time.Now()context.WithValue 塞进 ctx;在 OnEndFn 里取出来算 time.Since(start) 打印。其余三个时机(OnError / 两个流式)完全不用管——这正是 NewHandlerBuildercallbacks/handler_builder.go:109)的价值:直接实现 Handler 接口(callbacks/interface.go:85Handler = callbacks.Handler)要写全 5 个方法,用 Builder 只挑你关心的时机设,剩下的框架给空实现。这个 Handler 挂上后,图里每个节点自动被计时,业务代码一行没动。
L04

怎么挂上去

两种挂法:全局(对所有执行生效)或运行时(这一次执行生效):

// 方式一:全局挂(进程启动时,对之后所有调用生效)
callbacks.AppendGlobalHandlers(handler)

// 方式二:运行时挂(只对这一次 Invoke/Stream 生效)
r.Invoke(ctx, input, compose.WithCallbacks(handler))
读法:全局挂适合"所有请求都要的监控/日志";运行时挂适合"只想给这一次调用加临时观测"(比如调试某个请求)。两者可叠加。挂上后,框架执行每个节点时自动在 5 个时机调用你的 handler,你的业务/组件代码一行都不用改
这就是"无侵入"的含义 你写组件时压根不知道有没有人在监控它;运维想加监控时,也不用去改组件源码。两边解耦。第三方写的组件、Eino 内置的节点,照样能被你的 handler 观测到——因为埋点在框架的执行层,不在组件里。
L05

RunInfo:告诉你"现在在哪"

每次回调都带一个 RunInfo,告诉你当前是哪个节点、什么组件类型、叫什么名字

type RunInfo struct {
    Name      string   // 节点名
    Type      string   // 具体类型
    Component components.Component   // 组件类别(ChatModel/Tool/...)
}
读法:有了 RunInfo,你的 handler 就能区分"这次触发是模型节点还是工具节点",分别处理。比如只给 ChatModel 节点统计 token、只给 Tool 节点记录调用参数。链路追踪时用它给 span 命名。
L06

流式回调的坑(必须知道)

流式场景(OnEndWithStreamOutput)里,回调拿到的是一个 StreamReader(Day 03)。而流只能被消费一次——如果你的 handler 把流读完了,真正的业务下游就读不到了!

解法:Copy(Day 03 的 StreamReader.Copy)。框架/你在把流交给回调前会 Copy 一份,让回调读副本、业务读原本,互不影响。写流式 handler 时切记:要么用框架已 Copy 好的副本,要么自己 Copy——绝不能直接把业务的流读干。这是流式回调最常见的 bug 来源。另外流式 handler 读流是异步的,别在里面做重活阻塞主流程。
为什么流只能读一次? 回忆 Day 03:StreamReader 像一次性水管,水(chunk)流过去就没了。两个人都想喝,就得先分叉(Copy)成两根管。回调和业务都想看模型输出流,必须 Copy。这是流式一等公民设计带来的必然约束。

👶 小白:我在 OnEndWithStreamOutput 里就是想打个日志,把模型输出的流读出来看看内容——这能有什么问题?

👨‍🏫 老师:问题大了。接着"一次性水管"的比方——模型输出这根水管里的水(chunk)流过一次就没了。你在回调里把它读完打日志,等于把水全喝光,真正等着这股水的业务下游(前端、下一个节点)就一滴都接不到,界面上表现为"模型好像没输出"。

👨‍🏫 解法就是这节说的 Copy:把水管分叉成两根,回调读副本、业务读原本,各喝各的互不影响。实践上通常框架已经替你 Copy 好、把副本交给回调;你自己若要再分给别处,就自己再 Copy 一份——铁律是"绝不直接读业务那根流"。另外别在回调里对着流做重活(它是异步读的),拖慢主流程就本末倒置了:回调的本分是"只观测、不打扰"。

L07

典型用途

  • 日志:每个节点的输入输出打点,排查问题。
  • 指标监控:节点耗时、成功率、token 用量 → Prometheus。
  • 链路追踪:给每个节点开一个 span,串成一条 trace(OpenTelemetry/Langfuse)。
  • 成本核算:在模型节点的 OnEnd 里累加 token → 算这次请求花了多少钱(Day 01 成本经济学)。
🚨 错误驱动:没有回调(可观测)会出什么事故? 设想线上 Agent 突然变慢、偶尔答非所问,而你没挂任何回调:你看不到是哪个节点慢、模型到底吐了什么、工具返回了什么脏数据——只能靠"加 log 重新部署再复现",一轮几十分钟,半夜排障能查到天亮。挂了回调后,每个节点的耗时/输入/输出/错误全都记在案,出问题时打开这份 trace,几秒钟定位到"是 get_weather 工具超时导致模型反复重试"。可观测不是锦上添花,是线上事故时你唯一的眼睛。
eino-ext(Day 19)里就有现成的 callback handler,比如对接 Langfuse、APM 的,拿来即用。可观测性是生产级 Agent 的刚需——线上出问题时,这些埋点数据就是你的眼睛。
🗣️ 一句话复述 callbacks = 在框架执行层给每个节点装的"进出登记本",5 个时机(开始/成功/出错/流式进/流式出)自动触发你的 Handler,业务代码一行不改,就能拿到日志、耗时、追踪、成本——只观测、不打扰。
L08

今日小结 + 动手

🧠 今天你应该能回答

  • 回调解决什么?为什么叫 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
明天预告 · Day 17流式深水区——回到 Day 03 的 StreamReader,读它内部实现:Pipe 缓冲、Copy 的引用计数、多路合并、自动拼接。理解流式为什么能既快又安全。
← Day 15 ReAct Day 17 · 流式深入 →