您的位置:首页 > 手游攻略 > Agent 可观测性:Eino Callback 系统源码剖析(第59篇-E45)

Agent 可观测性:Eino Callback 系统源码剖析(第59篇-E45)

作者:互联网  时间: 2026-07-22 08:14:54  

读完这篇你会知道


为什么需要 Callback

一个 ReAct Agent 跑完一次,外部几乎看不见内部发生了什么:

Agent 可观测性:Eino Callback 系统源码拆解(第59篇-E45)

  • LLM 调用了几次?每次用了多少 Token?
  • 哪个工具被调用了?花了多长时间?
  • 哪个节点最慢,哪个 Retriever 召回质量最差?

没有这些数据,生产环境的问题排查和成本控制都是瞎猜。

Eino 的 Callback 系统是一个侵入性极低的观测钩子:不修改任何组件的业务逻辑,只在组件的生命周期关键点插入回调。


核心接口:Handler + 5 个 Timing

Handler 是所有观察者必须实现的接口:

 复制代码// 5个生命周期回调:
// OnStart / OnEnd / OnError / OnStartWithStreamInput / OnEndWithStreamOutput
type Handler interface {
    OnStart(ctx context.Context, info *RunInfo, input CallbackInput) context.Context
    OnEnd(ctx context.Context, info *RunInfo, output CallbackOutput) context.Context
    OnError(ctx context.Context, info *RunInfo, err error) context.Context    OnStartWithStreamInput(ctx context.Context, info *RunInfo,
        input *schema.StreamReader[CallbackInput]) context.Context
    OnEndWithStreamOutput(ctx context.Context, info *RunInfo,
        output *schema.StreamReader[CallbackOutput]) context.Context
}

每个方法的返回值是 context.Context。这个设计让同一个 Handler 的不同 timing 之间可以传递状态

 复制代码// OnStart 在 context 里存一个开始时间
func OnStart(ctx, info, input) context.Context {
    return context.WithValue(ctx, startKey{}, time.Now())
}// OnEnd 从 context 里取出开始时间,计算耗时
func OnEnd(ctx, info, output) context.Context {
    start := ctx.Value(startKey{}).(time.Time)
    log.Printf("[%s] 耗时: %v", info.Name, time.Since(start))
    return ctx
}

注意:这个 context 链只在同一个 Handler 的同一次调用内流动,不会跨 Handler。

RunInfo 描述了"是谁触发了这个回调":

 复制代码type RunInfo struct {
    Name      string              // 节点名(compose.WithNodeName 指定)
    Type      string              // 实现类型,如 "OpenAI"、"DeepSeek"
    Component components.Component // 组件种类:ChatModel/Tool/Retriever...
}

这三个字段是 Callback 逻辑里最重要的过滤条件:不同种类的组件要采集不同的指标。


5 个 Timing 的含义

Timing触发时机用于
OnStart组件开始处理前(非流式输入)记录入参、开始时间
OnEnd组件成功返回后(非流式输出)记录输出、耗时、Token 数
OnError组件返回错误时报错统计、告警
OnStartWithStreamInput组件接收流式输入较少用,Collect/Transform 场景
OnEndWithStreamOutput组件产出流式输出Token 流统计、流式日志

大多数使用场景只需要 OnStart + OnEnd + OnError


注册 Handler:两种方式

方式一:全局 Handler(进程级)

 复制代码// 在 main 或 TestMain 里调用,一次配置全部生效
// 注意:不是线程安全的,不能在 graph 运行时调用
callbacks.AppendGlobalHandlers(myTracingHandler, myMetricsHandler)

全局 Handler 对所有组件、所有 Graph 调用都生效。适合全平台的 OTel 追踪、Token 计量。

方式二:按调用注册

 复制代码// 只对这一次 Invoke 生效
result, err := runner.Invoke(ctx, input,
    compose.WithCallbacks(myDebugHandler),
)

这种方式适合开发调试或者租户级别的独立观测,不影响其他租户。


HandlerBuilder:最快写一个 Handler

不用实现全部 5 个方法,只订阅你关心的 timing:

 复制代码tokenCounter := callbacks.NewHandlerBuilder().
    OnEndFn(func(ctx context.Context, info *callbacks.RunInfo, output callbacks.CallbackOutput) context.Context {
        // 只处理 ChatModel 的输出,用 model.ConvCallbackOutput 安全类型转换
        mo := model.ConvCallbackOutput(output)
        if mo != nil && mo.TokenUsage != nil {
            metrics.Add("llm.tokens.total", float64(mo.TokenUsage.TotalTokens), map[string]string{
                "model": info.Name,
                "type":  info.Type,
            })
        }
        return ctx
    }).
    Build()

HandlerBuilder 内部实现了 TimingChecker——只注册了 OnEnd 的 Handler,在其他 timing 时框架会直接跳过,不分配 goroutine 也不复制 stream:

 复制代码func (hb *handlerImpl) Needed(_ context.Context, _ *RunInfo, timing CallbackTiming) bool {
    switch timing {
    case TimingOnEnd:
        return hb.onEndFn != nil   // 只注册了 OnEnd → 其他 timing 返回 false
    case TimingOnStart:
        return hb.onStartFn != nil // 没注册 → false,跳过
    // ...
    }
}

HandlerHelper:按组件类型分发

问题:CallbackInput / CallbackOutputinterface{},你在 OnEnd 里拿到的可能是 ChatModel 的输出,也可能是 Tool 的输出,需要手动判断类型。

HandlerHelper 把这个分发逻辑封装好了:

 复制代码helper := utils_callbacks.NewHandlerHelper().
    ChatModel(&utils_callbacks.ModelCallbackHandler{
        OnEnd: func(ctx context.Context, info *callbacks.RunInfo, output *model.CallbackOutput) context.Context {
            // 这里拿到的已经是强类型的 *model.CallbackOutput
            if output.TokenUsage != nil {
                log.Printf("[model:%s] tokens=%d", info.Name, output.TokenUsage.TotalTokens)
            }
            return ctx
        },
        OnEndWithStreamOutput: func(ctx context.Context, info *callbacks.RunInfo,
            output *schema.StreamReader[*model.CallbackOutput]) context.Context {
            // 流式输出:必须关闭 StreamReader,否则 goroutine 泄漏
            go func() {
                defer output.Close()
                var totalTokens int
                for chunk, err := output.Recv(); err == nil; chunk, err = output.Recv() {
                    if chunk.TokenUsage != nil {
                        totalTokens += chunk.TokenUsage.TotalTokens
                    }
                }
                log.Printf("[model:%s] stream total tokens=%d", info.Name, totalTokens)
            }()
            return ctx
        },
    }).
    Tool(&utils_callbacks.ToolCallbackHandler{
        OnStart: func(ctx context.Context, info *callbacks.RunInfo, input *tool.CallbackInput) context.Context {
            log.Printf("[tool:%s] called with: %s", info.Name, input.ArgumentsInJSON)
            return context.WithValue(ctx, toolStartKey{}, time.Now())
        },
        OnEnd: func(ctx context.Context, info *callbacks.RunInfo, output *tool.CallbackOutput) context.Context {
            start := ctx.Value(toolStartKey{}).(time.Time)
            log.Printf("[tool:%s] took=%v result=%s", info.Name, time.Since(start), output.Content)
            return ctx
        },
    }).
    Handler()callbacks.AppendGlobalHandlers(helper)

HandlerHelper 支持的组件类型:Prompt / ChatModel / Embedding / Indexer / Retriever / Loader / Transformer / Tool / ToolsNode / Agent + 对应的 Agentic 变体。


Callback 是怎么注入到 Graph 节点的

回到 E42 讲过的 taskManager.execute,它在每个节点执行前都会调用 initNodeCallbacks

 复制代码// compose/graph_manager.go
func (t *taskManager) execute(currentTask *task) {
    // 注入:把节点的 RunInfo + 本次 WithCallbacks 的 Handler 注入 context
    ctx := initNodeCallbacks(
        currentTask.ctx,
        currentTask.nodeKey,
        currentTask.call.action.nodeInfo,
        currentTask.call.action.meta,
        t.opts...,
    )
    // 真正执行节点
    task.output, task.err = runWrapper(ctx, call.action, ...)
}

initNodeCallbacks 的逻辑:

 复制代码// compose/utils.go
func initNodeCallbacks(ctx context.Context, key string, info *nodeInfo, meta *executorMeta, opts ...Option) context.Context {
    // 1. 从 meta 提取 RunInfo.Component、RunInfo.Type(组件种类 + 实现类型)
    ri := &callbacks.RunInfo{}
    if meta != nil {
        ri.Component = meta.component
        ri.Type = meta.componentImplType  // e.g. "OpenAI" / "DeepSeek"
    }
    // 2. 从 info 提取 RunInfo.Name(节点名,WithNodeName 指定的)
    if info != nil {
        ri.Name = info.name
    }    // 3. 收集本次 Invoke 传入的 WithCallbacks(handler)(仅路径匹配的)
    var cbs []callbacks.Handler
    for i := range opts {
        if len(opts[i].paths) > 0 {
            for _, k := range opts[i].paths {
                if k.path[0] == key {
                    cbs = append(cbs, opts[i].handler...)
                }
            }
        }
    }    // 4. 把全局 Handler + 路径匹配的 Handler 注入到 context
    if len(cbs) == 0 {
        return icb.ReuseHandlers(ctx, ri)   // 只用全局 Handler
    }
    return icb.AppendHandlers(ctx, ri, cbs...) // 全局 + 额外 Handler
}

这就是 Callback 的注入路径:Graph 执行前每个节点都会调用这个函数,把 RunInfoHandler 写进 context,节点内部的组件实现(ChatModel/Tool)调用 callbacks.OnStart(ctx, ...) 时就能找到它们。


流式回调的陷阱:必须关闭 StreamReader

当 ChatModel 流式输出时,OnEndWithStreamOutput 收到的是一个 *schema.StreamReader 副本(框架自动 stream.Copy)。

 复制代码// 错误做法:直接在 OnEndWithStreamOutput 同步消费
OnEndWithStreamOutput: func(ctx, info, output) context.Context {
    for chunk, _ := output.Recv() {  // 会阻塞,因为流还在传输
        // ...
    }
    // 忘记 output.Close() → goroutine 泄漏
    return ctx
},

正确做法:异步消费 + 一定要 Close

 复制代码OnEndWithStreamOutput: func(ctx, info, output) context.Context {
    go func() {
        defer output.Close() // 关闭这个副本,否则底层资源不释放
        for {
            chunk, err := output.Recv()
            if err != nil { break }
            // 处理 chunk...
        }
    }()
    return ctx  // 立即返回,不阻塞主流程
},

框架注释说得很清楚:


实战:用 Callback 接入 OTel Span

 复制代码import (
    "go.opentelemetry.io/otel"
    "go.opentelemetry.io/otel/attribute"
)var tracer = otel.Tracer("deepflux/agent")otlHandler := callbacks.NewHandlerBuilder().
    OnStartFn(func(ctx context.Context, info *callbacks.RunInfo, input callbacks.CallbackInput) context.Context {
        // 每个节点开始时创建一个 OTel Span
        ctx, span := tracer.Start(ctx, info.Name,
            trace.WithAttributes(
                attribute.String("eino.component", string(info.Component)),
                attribute.String("eino.type", info.Type),
            ))
        // 把 span 存进 context,供 OnEnd 结束它
        return context.WithValue(ctx, spanKey{}, span)
    }).
    OnEndFn(func(ctx context.Context, info *callbacks.RunInfo, output callbacks.CallbackOutput) context.Context {
        if span, ok := ctx.Value(spanKey{}).(trace.Span); ok {
            // Token 信息作为 Span attribute
            if mo := model.ConvCallbackOutput(output); mo != nil && mo.TokenUsage != nil {
                span.SetAttributes(
                    attribute.Int("llm.prompt_tokens", mo.TokenUsage.PromptTokens),
                    attribute.Int("llm.completion_tokens", mo.TokenUsage.CompletionTokens),
                )
            }
            span.End()
        }
        return ctx
    }).
    OnErrorFn(func(ctx context.Context, info *callbacks.RunInfo, err error) context.Context {
        if span, ok := ctx.Value(spanKey{}).(trace.Span); ok {
            span.RecordError(err)
            span.SetStatus(codes.Error, err.Error())
            span.End()
        }
        return ctx
    }).
    Build()callbacks.AppendGlobalHandlers(otlHandler)

这样每个节点的执行都会在 Jaeger/Tempo 里有对应的 Span,自动组成 Trace 树。


小结

Eino Callback 系统的设计核心是最小侵入 + 类型安全

机制作用
Handler 接口 + 5 个 Timing覆盖组件全生命周期,统一契约
RunInfo 三字段标识"谁在运行",用于 Metric label / Span attribute
TimingChecker.Needed避免没注册的 timing 产生不必要的 stream 复制
HandlerBuilder只订阅关心的 timing,其余跳过
HandlerHelper按组件类型强类型分发,无需手写 type switch
initNodeCallbacksGraph 节点执行时注入 RunInfo + Handler,不需要组件感知 Graph
流式副本 + Close流式 Callback 的资源管理约定,漏 Close = goroutine 泄漏

在 DeepFlux 的 server/internal/observability 模块里,就是用这套机制把每次 Agent 调用的 Token 消耗、工具调用延迟、Retriever 召回延迟,以 OTel Span 的形式发送给 Tempo,再通过 Grafana 展示出来,完整覆盖从 HTTP 入口到 LLM 调用的链路追踪。


代码来源:eino/callbacks/interface.go · eino/callbacks/handler_builder.go · eino/utils/callbacks/template.go · eino/compose/utils.go

最新游戏

更多

Copyright©2010-2019. All rights reserved | 波波三国游戏官网|[email protected]

备案编号:湘ICP备2022015115号-4