KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
06 · 流式与重订阅:进展怎么实时取回来 — keel 龙骨
## 现场:要么刷屏,要么丢帧
现场:要么刷屏,要么丢帧
诊断任务要跑二十分钟。第一版实现用轮询:每 5 秒一次 GetTask,二十分钟就是 240 次请求,其中 238 次拿到的状态没变。用户那边则是一个转了二十分钟的圈,什么也看不到。
改成 SSE 流式之后体验好了:状态一变就推一帧。但上线两周出了新问题——用户在会议室网络切换了三秒,SSE 连接断了,重连之后客户端从零开始收,前面已经推过的进度帧全部重复,而断线那三秒里发生的状态跃迁则彻底丢了。
两个版本的问题其实是同一个:「取回进展」这件事没有可靠的断点续传。
直觉模型:同一批对象,两种取回方式
先把一件事说清楚——流式不改变内容,只改变送达方式。
| 轮询 | 流式(SSE) | |
|---|---|---|
| 谁主动 | 客户端反复问 | 服务端持续推 |
| 拿到的对象 | Task |
Task 加一串更新事件 |
| 断线后 | 下次轮询自然补齐 | 需要 SubscribeToTask 补回 |
| 适合 | 短任务、低频 | 长任务、要实时反馈 |
本项目里两者共用同一段执行逻辑,emit 回调只是决定「这一帧要不要顺手记下来」——这就是「流式只是取回方式不同」的代码证据。
精确定义:靠成员名判别,而不是 kind
v1.0 的流式事件有两种:
| 事件 | 判别用的成员名 | 携带 |
|---|---|---|
TaskStatusUpdateEvent |
statusUpdate |
taskId contextId status |
TaskArtifactUpdateEvent |
artifactUpdate |
taskId contextId artifact index |
两处破坏性变更要记住(来源:a2a-protocol.org What's New in v1.0,检索于 2026-10-05):
- 去掉了
kind判别器,kind: "status-update"改成了外层成员名statusUpdate。 - 去掉了
final布尔。任务是否到达终态,由协议绑定的流关闭机制表示,不再由某一帧的字段表示。
第二条尤其重要:如果你的客户端还在等 final: true,它永远等不到——因为 v1.0 根本不发这个字段。正确的做法是看流什么时候关闭。
另外两点规范澄清:允许对同一任务建立多个并发订阅,且每个订阅收到相同顺序的事件;SubscribeToTask 用于断线后重连补帧。
全链路图
flowchart TD
subgraph CLIENT["Client Agent"]
S1["① SendStreamingMessage"] --> OPEN["② SSE 连接建立"]
OPEN --> LOOP["③ 逐帧处理"]
LOOP --> DROP["④ 连接中断"]
DROP --> RESUB["⑤ SubscribeToTask 重订阅"]
end
subgraph SERVER["Remote Agent"]
EXEC["执行器推进任务"]
EMIT["产出 statusUpdate 或 artifactUpdate"]
CLOSE["任务到终态 关闭流"]
end
S1 --> EXEC
EXEC --> EMIT
EMIT -->|通过 SSE data 行| LOOP
EXEC --> CLOSE
CLOSE -->|流关闭即终态信号| LOOP
RESUB --> BACKFILL["⑥ 补回当前状态与已有产物"]
BACKFILL --> LOOP
S1 -->|服务端未声明 streaming| UNSUP["错误 UNSUPPORTED_OPERATION"]
UNSUP -.-> CLIENT
DROP -.->|未重订阅 丢帧| GAP["断线期间的状态跃迁永久丢失"]
一次完整运行
python courses/foundation/a2a-protocol-engineering/course/project/examples/05_streaming.py
实跑输出:
A. SendStreamingMessage 产出的帧:
1. task 任务建立 state=TASK_STATE_SUBMITTED
2. statusUpdate state=TASK_STATE_WORKING
3. statusUpdate state=TASK_STATE_COMPLETED
4. artifactUpdate artifact=artifact-01 index=0
B. 判别方式: 靠顶层成员名,不靠 kind 字段
frame 里有没有 kind: False
C. 断线后重订阅: 补回 2 帧 | stream closed = True
D. 服务端没声明 streaming 却调用: UNSUPPORTED_OPERATION
逐段读:
- A 是完整的一串帧:先建任务,再两次状态跃迁,最后产物。注意第 1 帧是
task,不是statusUpdate——建立任务本身也是一帧。 - B 是判别方式:整串帧里没有任何
kind字段。客户端必须用"statusUpdate" in frame这种写法。 - C 是重订阅:补回 2 帧(当前状态 + 已有产物),并且
closed=True表示这个任务已经终态、流可以关了。 - D 是能力未声明时的拒绝:服务端卡面上没写
streaming: true,调用就被拒。这和 MCP 里的能力协商是同一个道理——先声明,再调用。
失败注入
注入 A:把 final 当终态
v0.3 迁移过来的老代码常常这样写:
if event.get("final"):
break
在 v1.0 上这个条件永远为假(字段不存在),客户端会一直等到连接超时。修复是改为监听流关闭:
for frame in stream: # 流结束即终态
handle(frame)
settle() # 循环结束 = 服务端关闭 = 任务到终态
注入 B:断线后从零开始
模拟连接中断(丢掉已收到的帧),然后不调 SubscribeToTask,直接重新 SendStreamingMessage 再跑一次。结果是任务被重跑一遍——对有副作用的任务,这是实打实的重复执行。
正确顺序是:断线 → SubscribeToTask 补帧 → 只有在补不回来(任务不存在或已终态)时才考虑重发。
生产边界
| 教学实现 | 生产替换 |
|---|---|
| 一次性返回帧列表 | 真实 SSE(text/event-stream),每条 data: 一个 JSON-RPC 对象 |
closed 布尔 |
HTTP 层的流关闭事件 |
| 无重连退避 | 指数退避重连 + Last-Event-ID 或 SubscribeToTask |
| 无心跳 | 定期心跳帧,防止代理与负载均衡器掐断空闲连接 |
| 无并发控制 | 限制单任务订阅数、按客户端限流 |
练习与验收
练习(有可观察结果):写一个 consume(frames) 函数,输入一串帧,返回「最终状态」和「产物列表」。要求:① 靠成员名判别而不是 kind;② 同一个 artifactId 的多个 artifactUpdate 按 index 去重,只保留最后一次;③ 遇到流关闭(本实现里是 closed=True)才认定终态。
验收标准:喂给它一串包含重复 artifactUpdate(相同 artifactId 与 index)的帧,断言产物只出现一次。如果你的实现出现重复产物,说明你把「传输层的补帧」当成了「业务层的新产物」——这是重订阅场景里最常见的重复来源。
本章检查点
- 现场两个版本(轮询刷屏、流式丢帧)分别缺什么?
- 为什么 v1.0 去掉
final之后,客户端必须改监听流关闭?不改会怎样? - 断线之后,重订阅和重新发送,哪一个才是正确的第一步?为什么另一个危险?
现在能解释什么
你现在能解释流式只是取回方式的变化:同一段执行逻辑,既可以被轮询消费,也可以被 SSE 消费。你也知道 v1.0 用成员名判别事件、用流关闭表示终态,以及断线后必须走 SubscribeToTask 补帧而不是重跑任务。下一章解决最后一种取回方式:客户端根本不在线时怎么办——推送通知。