vault backup: 2026-04-29 20:17:22
This commit is contained in:
@@ -324,6 +324,10 @@ for {
|
||||
> - **不要忽略 `event.Err`**——Agent 内部可能静默失败(如工具执行超时)
|
||||
> - **注意 goroutine 安全**——多个消费者同时读取同一个 AsyncIterator 是不安全的
|
||||
|
||||
> [!question] 深入理解事件流消费模式?
|
||||
> 通过 Claude Code Agent 事件流消费的类比加深理解。
|
||||
> -> 参考 [[chapter_02_chatmodelagent_runner_agentevent/async_iterator_consumption|AsyncIterator:事件流的消费方式]]
|
||||
|
||||
## 多轮对话的实现
|
||||
|
||||
本章实现的是简单的多轮对话:用户输入 → 模型回复 → 用户继续输入 → ...
|
||||
|
||||
+94
@@ -0,0 +1,94 @@
|
||||
---
|
||||
tags: [eino, asynciterator, go, streaming]
|
||||
create time: 2026-04-29 15:30
|
||||
---
|
||||
|
||||
# AsyncIterator:事件流的消费方式
|
||||
|
||||
## 概述
|
||||
|
||||
理解 Eino 的 `AsyncIterator` 最简单的方式——**看 Claude Code 是怎么消费自身事件的**。两者的消费过程完全一致。
|
||||
|
||||
---
|
||||
|
||||
## Claude Code 的事件消费过程
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A["runner.Run() 获取事件流"] --> B["循环 events.Next()"]
|
||||
B --> C{"事件类型?"}
|
||||
C -->|"thinking"| D["显示思考过程"]
|
||||
C -->|"text_delta"| E["追加文本到终端"]
|
||||
C -->|"tool_use"| F["高亮工具调用"]
|
||||
C -->|"error"| G["记录错误,终止"]
|
||||
C -->|"interrupt"| H["用户取消,优雅退出"]
|
||||
C -->|"无更多事件"| I["closeConnection()"]
|
||||
|
||||
D --> B
|
||||
E --> B
|
||||
F --> B
|
||||
style A fill:#e1f5fe
|
||||
style I fill:#ffebee
|
||||
```
|
||||
|
||||
核心过程只有一句:**获取流 → Next() 逐个拉取 → 按类型处理 → Close 释放。**
|
||||
|
||||
---
|
||||
|
||||
## Eino 的等价过程
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A["runner.Run() 获取迭代器"] --> B["循环 events.Next()"]
|
||||
B --> C{"event.Err?"}
|
||||
C -->|是| G["记录错误,退出循环"]
|
||||
C -->|否| D{"event.Output?"}
|
||||
D -->|是| E["打印内容给用户"]
|
||||
D -->|否| F{"event.Action?"}
|
||||
F -->|是| J["处理控制动作<br/>本章节用不到"]
|
||||
F -->|否| B
|
||||
E --> B
|
||||
G --> H["events.Close() 释放资源"]
|
||||
|
||||
style A fill:#e1f5fe
|
||||
style H fill:#ffebee
|
||||
```
|
||||
|
||||
同样四个字阶段:获取流 → Next() 逐个拉取 → 按类型处理 → Close 释放。
|
||||
|
||||
---
|
||||
|
||||
## 两者对照
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
subgraph Claude Code
|
||||
A["thinking"] --> T1["显示思考过程"]
|
||||
B["text_delta"] --> T2["追加文本输出"]
|
||||
C["tool_use"] --> T3["高亮工具调用"]
|
||||
D["error"] --> T4["错误终止"]
|
||||
E["interrupt"] --> T5["用户取消"]
|
||||
end
|
||||
|
||||
subgraph Eino ADK
|
||||
A1["—"] --> S1["无此概念"]
|
||||
B1["event.Output"] --> S2["打印消息内容"]
|
||||
C1["event.Action"] --> S3["内部调度信号"]
|
||||
D1["event.Err"] --> S4["错误退出"]
|
||||
E1["ctx 被 cancel"] --> S5["用户取消"]
|
||||
end
|
||||
|
||||
B -.等价.-> B1
|
||||
C -.等价.-> C1
|
||||
D -.等价.-> D1
|
||||
E -.等价.-> E1
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 注意事项
|
||||
|
||||
- 始终检查 `Err` 字段——错误通常是静默发生的
|
||||
- `Next()` 返回的 `ok` 为 false 时表示流已结束,不应继续调用
|
||||
- 每次 `Run()` 创建的迭代器只能消费一次,不能复用
|
||||
- 不手动调用 `Close()` 会导致资源泄漏
|
||||
Reference in New Issue
Block a user