0%

生产级 Agent 架构:以 Eino 搭建 Chatbot 为例

1. 概览

一个 Chatbot 要上线,关键不在模型多强,在模型外面套一圈确定性的关卡。

  • 违规请求在进模型之前就被拦掉,省 token 也省风险。
  • 召回来的知识先过重排再喂给模型,低分的直接砍掉。
  • 工具调用套上中间件,限流、清洗、重试都在这一层做,下游不会被模型不可靠的参数打崩,出错能自愈。
  • 输出过审才出去、过审才写回记忆,某一轮的错结论不会污染之后每一轮。
  • 全程挂 Callback,出问题能定位到是哪个环节慢、哪个环节多花了 token。

思路是:模型能出错的每一处,都放一道确定性的关卡接着,模型只管认知,工程交给外面的关卡。

1
2
3
4
5
6
7
8
9
10
11
12
┌────────────────────────────────────────┐
│ 外层 compose.Graph:确定性的关卡 │
│ │
│ ┌──────────────────────────────────┐ │
│ │ 内核 adk ReAct:模型 ↔ 工具循环 │ │
│ │ │ │
│ │ ┌────────────────────────────┐ │ │
│ │ │ 工具护栏(中间件链) │ │ │
│ │ │ 限流 → 清洗 → 校验重试 │ │ │
│ │ └────────────────────────────┘ │ │
│ └──────────────────────────────────┘ │
└────────────────────────────────────────┘

关卡之下,还有几类横贯全程的底座:中间件改行为,Callback 做观测,State 隔离状态,流扛背压。结构看下面的全景图,落地看第 4 节的伪代码。

2. 详细架构图

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
┌ Callback 探针(贯穿全图:每个节点、每次模型/工具调用都挂)

│ [ HTTP / SSE 请求:Query + SessionID + UserID ]
│ │
│ ▼

│ ═══════════════ 接入与调度层 · compose.Graph ═══════════════

│ ┌──────────────────────────────────┐
│ │ ① 输入风控 │ ──违规──▶ END
│ │ 违规请求不进模型 │ (不花模型调用)
│ └─────────────────┬────────────────┘
│ │
│ │ 放行
│ ┌──────────────┴──────────────┐
│ ▼ 并发 ▼ 并发
│ ┌──────────────────────┐ ┌──────────────────────┐
│ │ ② 记忆加载 │ │ ③ 知识召回 │
│ │ 会话历史 Redis │ │ 向量 Dense │
│ │ 用户画像 │ │ 关键词 BM25 │
│ └───────────┬──────────┘ └───────────┬──────────┘
│ └──────────────┬──────────────┘
│ ▼
│ ┌──────────────────────────────────┐
│ │ ④ 重排 + 话题判别 │
│ │ Cross-Encoder 二次打分 │
│ │ 砍掉低于阈值的切片 │
│ │ 淘汰跨话题的失效记忆 │
│ └─────────────────┬────────────────┘
│ ▼
│ ┌──────────────────────────────────┐
│ │ ⑤ 拼装 + Token 预算 │
│ │ system / 知识 / 历史 / 输出 │
│ │ 按比例压进上下文窗口 │
│ └─────────────────┬────────────────┘

│ ═══════════════ 核心推理层 · adk ReAct ═══════════════

│ ▼
│ ┌──────────────────────────────────┐
│ │ ChatModel(思考 / 规划) │
│ │ ├─ 触发 ToolCall ──▶ 工具洋葱 │
│ │ └─ 输出最终回复 ──▶ 退出循环 │
│ └─────────────────┬────────────────┘
│ ▼
│ ┌──────────────────────────────────┐
│ │ 工具洋葱(中间件链) │
│ │ M1 并发限流(信号量) │
│ │ M2 参数清洗(剥 markdown) │
│ │ M3 校验重试(指数退避) │
│ │ M4 真实工具(业务 / 检索) │
│ └─────────────────┬────────────────┘
│ │ 结果喂回 ChatModel,循环继续
│ │ 退出循环
│ ▼

│ ═══════════════ 退出与沉淀层 · compose.Graph ═══════════════

│ ┌──────────────────────────────────┐
│ │ ⑥ 输出风控 + 脱敏 │
│ │ 命中 → 换兜底话术 │
│ └─────────────────┬────────────────┘
│ ▼
│ ┌──────────────────────────────────┐
│ │ ⑦ 记忆回写 │
│ │ 过审才落库,异步 │
│ └─────────────────┬────────────────┘
│ ▼
│ [ SSE 流式返回 ]
└ TraceID 树状传播 · token 计数 · 耗时统计 · 告警

图里画得出的是节点和连线,画不出的还有两类底座:

1
2
State   隔离状态:编译期只读单例,运行期每请求一份,没有跨请求的锁
流 扛背压:有界 channel,下游慢就阻塞上游(直连流式下能传到 TCP 零窗口)

3. 一次请求的时序

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
Client       Graph(关卡)        Memory/Rerank     ReAct 内核      工具洋葱
│ │ │ │ │
│─ 请求 ──────▶│ │ │ │
│ │─ ① 输入风控 ──────────────────────────────────────▶ 违规则短路
│ │─ ② 记忆加载 ─────▶│ │ │
│ │─ ③ 知识召回 ─────▶│ │ │
│ │◀─ 历史+画像+切片 ─┘ │ │
│ │─ ④ 重排砍低分 ───▶│ │ │
│ │─ ⑤ 拼装+预算 ────────────────────▶│ │
│ │ │ │── ToolCall ─▶│
│ │ │ │ │─ M1 限流
│ │ │ │ │─ M2 清洗
│ │ │ │ │─ M3 校验重试
│ │ │ │ │─ M4 真实调用
│ │ │ │◀─ 结果 ──────┘
│ │ │ │─ 继续推,出最终回答
│ │◀─ 最终回答 ────────────────────────┘ │
│ │─ ⑥ 输出风控 ────────────────────────────────────────
│ │─ ⑦ 记忆回写(异步) ─────────────────────────────────
│◀─ SSE 流 ────┘ │ │ │

4. 串起来:完整伪代码

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
// ── 状态结构:在七个节点之间流转 ────────────────────────────────
type State struct {
Query string
SessionID string
UserID string

History []*schema.Message // 短时记忆
Profile UserProfile // 用户画像
Chunks []Document // 召回的知识切片

Blocked bool // 入口风控标记
Final *schema.Message // 最终回复
}

// ── 一、可观测回调:只能看,不能改 ────────────────────────────
// 跟着调用传进去(compose.WithCallbacks),或全局 AppendGlobalHandlers。
// 父子 span 靠图拓扑自动维护,定位慢点靠 span 树,不靠日志时间戳。
// 别在回调里改 Input/Output:下游节点和所有 handler 共享同一指针,改了会数据竞争。
trace := callbacks.NewHandlerBuilder().
OnStartFn(func(ctx context.Context, info *callbacks.RunInfo, _ callbacks.CallbackInput) context.Context {
return openSpan(ctx, info.Name) // 开 span,透传 TraceID
}).
OnEndFn(func(ctx context.Context, info *callbacks.RunInfo, _ callbacks.CallbackOutput) context.Context {
closeSpan(ctx) // 记耗时、记 token
return ctx
}).
OnErrorFn(func(ctx context.Context, info *callbacks.RunInfo, err error) context.Context {
recordError(ctx, info.Name, err) // 记异常堆栈
return ctx
}).
Build()

// ── 二、工具洋葱:中间件链,M1 → M4 ───────────────────────────
// 中间件签名是 func(endpoint) endpoint,一串套起来就是洋葱。
// Invokable 中间件只作用于实现了 InvokableTool 的工具,Streamable 工具不生效。
// 底层错误在这里截胡成诊断信息喂回模型,让模型自愈,不往上抛。
func retryAndValidate(next compose.InvokableToolEndpoint) compose.InvokableToolEndpoint {
return func(ctx context.Context, in *compose.ToolInput) (*compose.ToolOutput, error) {
var out *compose.ToolOutput
err := retry(ctx, 3, func() error { // M3:指数退避重试
o, e := next(ctx, in)
out = o
if e == nil && !valid(o.Result) { // 结果校验,不合法重试
return errInvalidResult
}
return e
})
return out, err
}
}

// ── 三、内核 ReAct agent ──────────────────────────────────────
// Instruction 是 system prompt,默认实现会拿会话变量做 FString 格式化,
// prompt 里带 {}(比如 JSON 示例)会报错,要么自定义 GenModelInput,要么别带花括号。
// 降级的 GetFailoverModel 里别 new 模型:实例提前建好,降级那一刻最需要快。
agent, _ := adk.NewChatModelAgent(ctx, &adk.ChatModelAgentConfig{
Name: "assistant",
Instruction: systemPrompt,
Model: chatModel,
MaxIterations: 10, // 卡死最大轮数,防止模型在循环里出不来
ModelRetryConfig: &adk.TypedModelRetryConfig[*schema.Message]{
MaxRetries: 3,
ShouldRetry: shouldRetry, // 判哪些错误值得重试
// BackoffFunc 默认指数退避 100ms → 10s + jitter
},
ModelFailoverConfig: &adk.ModelFailoverConfig[*schema.Message]{
MaxRetries: 1,
ShouldFailover: shouldFailover, // 判要不要换模型
GetFailoverModel: getFailoverModel, // 换到哪个模型
},
ToolsConfig: adk.ToolsConfig{
ToolsNodeConfig: compose.ToolsNodeConfig{
Tools: []tool.BaseTool{orderQuery, refundTool, searchDoc},
ToolCallMiddlewares: []compose.ToolMiddleware{
{Invokable: concurrencyLimiter(4)}, // M1 信号量限流,防打崩下游
{Invokable: argsSanitizer}, // M2 剥 markdown、修参数
{Invokable: retryAndValidate}, // M3 校验重试
},
UnknownToolsHandler: hallucinationFallback, // 模型幻觉出不存在的工具时兜住
},
},
Handlers: []adk.ChatModelAgentMiddleware{auditHandler}, // 改状态的中间件,和回调是两回事
})

// ── 四、外层图:确定性关卡(括号内是节点 key)─────────────────
// ① 输入风控(input_guard):命中违规就短路,不花一次模型调用
// ② 记忆加载 + ③ 知识召回(recall):一个节点并发拉会话历史 + 画像 + 向量/BM25
// ④ 重排(rerank):cross-encoder 二次打分,砍低分切片
// 中文用 bge-reranker-v2-m3,英文 bge-reranker-large,API 有 Cohere / Jina;只对 top-N 重排
// ⑤ 拼装 + 预算(assemble):system/知识/历史/输出空间按比例压进上下文窗口
// ⑥ 输出风控(output_guard):命中换兜底话术
// ⑦ 记忆回写(writeback):过审才落库,防模型输出污染下一轮记忆
// 内核 ReAct(react):是关卡保护的对象,不是关卡
g := compose.NewGraph[*State, *State]()

_ = g.AddLambdaNode("input_guard", compose.InvokableLambda(inputGuard))
_ = g.AddLambdaNode("recall", compose.InvokableLambda(recall))
_ = g.AddLambdaNode("rerank", compose.InvokableLambda(rerank))
_ = g.AddLambdaNode("assemble", compose.InvokableLambda(assemble))
_ = g.AddLambdaNode("react", compose.InvokableLambda(runReAct))
_ = g.AddLambdaNode("output_guard", compose.InvokableLambda(outputGuard))
_ = g.AddLambdaNode("writeback", compose.InvokableLambda(writeBack))

// 入口风控分支:命中跳 END,跳过后面所有节点。
// input_guard 这个 key 出现三次:AddLambdaNode 定义节点、AddEdge(START, ...) 接成入口、
// 这里的 AddBranch 决定出口。有分支的节点出口走 AddBranch,不再写 AddEdge。
g.AddBranch("input_guard", compose.NewGraphBranch(
func(ctx context.Context, s *State) (string, error) {
if s.Blocked {
return compose.END, nil
}
return "recall", nil
},
map[string]bool{"recall": true, compose.END: true},
))

g.AddEdge(compose.START, "input_guard")
g.AddEdge("recall", "rerank")
g.AddEdge("rerank", "assemble")
g.AddEdge("assemble", "react")
g.AddEdge("react", "output_guard")
g.AddEdge("output_guard", "writeback")
g.AddEdge("writeback", compose.END)

// 编译一次,进程内只读单例,被所有请求共享;运行期每请求独立 state,无跨请求锁。
// 注意 Compile 只收 GraphCompileOption,compose.WithCallbacks 是调用期选项,传不进这里,
// 要挂在 runnable.Invoke / Stream 上(见下面 SSE 那段)。
runnable, _ := g.Compile(ctx)

// 召回节点:记忆和知识并发拉,慢 IO 别串行。
func recall(ctx context.Context, s *State) (*State, error) {
g, ctx := errgroup.WithContext(ctx)
g.Go(func() error { s.History = loadSession(ctx, s.SessionID); return nil })
g.Go(func() error { s.Profile = loadProfile(ctx, s.UserID); return nil })
g.Go(func() error { s.Chunks = hybridSearch(ctx, s.Query); return nil })
if err := g.Wait(); err != nil {
return nil, err
}
return s, nil
}

// 内核节点:跑 agent 的 ReAct 循环,取最终回答。
func runReAct(ctx context.Context, s *State) (*State, error) {
msgs := assembleMessages(s) // system + 知识 + 历史 + 用户问题
// 不开流式:模型输出直接落在 Message 上。
// 开了流式(EnableStreaming: true)Message 是零值,正文在 MessageStream 里,
// 必须读到 !ok 再拼装,中途 return 会让发送协程永远悬在发送上。
iter := agent.Run(ctx, &adk.AgentInput{Messages: msgs})
for {
ev, ok := iter.Next()
if !ok {
break
}
if ev.Output == nil || ev.Output.MessageOutput == nil {
continue
}
m := ev.Output.MessageOutput
if m.Role != schema.Assistant {
continue // 工具结果也会走事件,别当成回答
}
s.Final = m.Message // 最后一条 assistant 消息是最终回答
}
return s, nil
}

// ── 五、SSE 服务端 ────────────────────────────────────────────
func handle(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
flusher, _ := w.(http.Flusher)

state := &State{
Query: r.URL.Query().Get("q"),
SessionID: r.Header.Get("X-Session-Id"),
UserID: r.Header.Get("X-User-Id"),
}

// trace 挂在这里:Compile 不接受 callbacks,调用期传才是生效的位置。
sr, err := runnable.Stream(r.Context(), state, compose.WithCallbacks(trace))
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}

// 注意这里拿到的是 StreamReader,不是 agent.Run 返回的 AsyncIterator:
// 前者的结束信号是 io.EOF,后者才是 Next() 的第二个返回值。
// 两种都要读到结束,中途 return 会让发送端 goroutine 悬在发送上。
defer sr.Close() // 即使读完也要关,注释里写明了 always close, even after io.EOF
for {
chunk, err := sr.Recv()
if errors.Is(err, io.EOF) {
break
}
if err != nil {
break // 流断了,别把半个响应当成功返回
}
fmt.Fprintf(w, "data: %s\n\n", chunk.Final.Content)
flusher.Flush()
}
// 这里 react 是 invoke 节点,Stream 返回整段;要 token 级打字机,把 react
// 换成 StreamableLambda,图会在流式和值节点之间自动插转换层。
}

5. 参考

  • cloudwego/eino · Graph、Branch、ToolCallMiddlewares、AgentMiddleware、callbacks、ADK 的源码,本文 API 名称以 main 分支为准
  • BAAI/bge-reranker-v2-m3 · 多语言 cross-encoder 重排模型
  • Cohere Rerank · 重排 API