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):

  1. 去掉了 kind 判别器,kind: "status-update" 改成了外层成员名 statusUpdate。
  2. 去掉了 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

逐段读:

  1. A 是完整的一串帧:先建任务,再两次状态跃迁,最后产物。注意第 1 帧是 task,不是 statusUpdate——建立任务本身也是一帧。
  2. B 是判别方式:整串帧里没有任何 kind 字段。客户端必须用 "statusUpdate" in frame 这种写法。
  3. C 是重订阅:补回 2 帧(当前状态 + 已有产物),并且 closed=True 表示这个任务已经终态、流可以关了。
  4. 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)的帧,断言产物只出现一次。如果你的实现出现重复产物,说明你把「传输层的补帧」当成了「业务层的新产物」——这是重订阅场景里最常见的重复来源。

本章检查点

现在能解释什么

你现在能解释流式只是取回方式的变化:同一段执行逻辑,既可以被轮询消费,也可以被 SSE 消费。你也知道 v1.0 用成员名判别事件、用流关闭表示终态,以及断线后必须走 SubscribeToTask 补帧而不是重跑任务。下一章解决最后一种取回方式:客户端根本不在线时怎么办——推送通知。

进入 keel 阅读