KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
04 · 三个真实场景生成器拆解 — keel 龙骨
前三章的语言知识,现在对到生产代码上。三个例子来自同一类真实系统——一个流式服务:模型逐 token 产出、中间层逐事件落库、SSE 逐帧推给浏览器。每个例子先给「它解决什么问题」,再逐行拆。示例都是通用可运行的结构,不绑定任何具体项目。
前三章的语言知识,现在对到生产代码上。三个例子来自同一类真实系统——一个流式服务:模型逐 token 产出、中间层逐事件落库、SSE 逐帧推给浏览器。每个例子先给「它解决什么问题」,再逐行拆。示例都是通用可运行的结构,不绑定任何具体项目。
例一:run_pipeline —— 事件泵(异步生成器 + 旁路观察)
一次流式运行的编排函数,典型写法:
async def run_pipeline(req, runner=None, hooks=None):
context = None
is_cancelled = False
alert_sent = False
async def send_alert(error_message): # 闭包:读外层最新状态
nonlocal alert_sent
if alert_sent:
return
await push_alert(error_message)
alert_sent = True
try:
context = await prepare_context(req) # 备料
if req.stream:
async for item in stream_events(context, req): # ← 消费上游生成器
if item.type == "cancelled":
is_cancelled = True # 旁路观察:记下取消
if item.type == "error":
mark_failed(context.run_id, item.error)
yield item # ← 原样往外交给下一层
else:
response = await run_once(context, req)
yield response
except Exception as exc:
mark_failed(context.run_id, str(exc))
await send_alert(str(exc))
raise
finally:
if is_cancelled:
mark_stopped(context.run_id)
await hooks.on_end() # 无论成败都触发
这个例子的三个知识点:
- 「消费上游 + 产出下游」的管道形态:
async for item in 上游+yield item,自己既是消费者又是生产者。事件链路里这种「中转层」往往有好几层,每层只加一点自己的逻辑(这里就是标记取消/错误)——每层都很薄,串起来才是完整管道; yield item在 try 块里:下游消费者出异常时,异常会从 yield 处抛进本函数,被 except 接住走兜底——这是 02 章 throw 机制的自动版(消费者不礼貌退出时,生产者能感知);- finally 收尾:无论管道怎么断,收尾钩子都会执行——对应 02 章 finally 清理的工程化用法。
例二:sse_generator —— 无限循环的心跳泵(本课最值得抄写的代码)
SSE 订阅端点的生成器,典型生产写法:
async def sse_generator(request, stream_key, start_offset):
current_offset = start_offset or "0-0"
last_event_time = time.time()
last_heartbeat_time = time.time()
heartbeat_interval = 30 # 心跳间隔(秒)
max_idle_seconds = 30 * 60 # 空闲熔断(30 分钟)
while True: # ← 无限循环:SSE 是长连接
if await request.is_disconnected(): # ① 消费者没了 → 退出
break
if time.time() - last_event_time > max_idle_seconds: # ② 空闲熔断
break
records, current_offset, done = await read_batch(
stream_key, current_offset, block_ms=5000) # ③ 阻塞读(await,让出控制权)
if not records:
if time.time() - last_heartbeat_time >= heartbeat_interval:
yield f"data: {json.dumps(heartbeat_payload)}\n\n" # ④ 产出心跳
continue # ← 没数据:continue 回到循环头
last_event_time = time.time() # 有真实事件,重置计时
last_heartbeat_time = time.time()
for record in records: # ⑤ 有数据:逐条产出
payload = record.payload | {"offset": record.offset} # 每帧带游标
yield f"data: {json.dumps(payload)}\n\n"
if done: # ⑥ 终态事件 → 正常结束
break
read_batch是一个「按游标批量取记录」的抽象(可以是 Redis Stream、数据库变更表、消息队列,随你)——它只要返回(records, 下一游标, 是否终态)即可,本例不关心底层。
用前三章的语言逐个解释设计:
| 代码 | 语言机制 | 为什么 |
|---|---|---|
while True + 多个 break |
生成器可以表达无限序列 | SSE 长连接不知道何时结束 |
yield f"data: ..." |
每帧产出即推送 | 消费方(StreamingResponse)拿到一帧立刻写 socket |
await read_batch(③) |
yield 之间可以 await(03 章核心) | 等数据的 5 秒里让出事件循环,服务器还能服务别人 |
is_disconnected() break(①) |
停止消费即停止生产 | 断开连接不需要通知生产者,下一轮自查即退 |
| 心跳 yield(④) | 无数据时主动产出 | 保活中间代理,同时证明「我还活着、我还能 yield」 |
| 每帧带 offset(⑤) | 续传游标内建在数据里 | 客户端断线重连时带上 offset,就能从断点继续——不需要额外握手协议 |
done break(⑥) |
有限序列的自然终点 | 终态事件后连接正常关闭 |
这个函数是「异步生成器 = 长连接生命周期管理」的完整范本:启动、运行(await 读 + yield 推)、退出(断开/超时/终态三路 break)、清理。学会它,你自己写 WebSocket/长轮询/推送服务都是同一副骨架。
例三:消费端 —— async for 的另一半
任务 Worker 里消费事件流的代码:
event_iterator = await build_event_source(run_id, request) # 返回生成器对象,此刻任务一步未跑
is_stopped = False
async for item in event_iterator:
if await should_stop(run_id): # 消费间隙做检查
cancel_current() # → 协作式取消
if item.type == "event":
await publish_event(run_id, item) # 写存储 / 落库
if item.type == "cancelled":
is_stopped = True
消费者视角的三个知识点:
await build_event_source(...)返回的是生成器对象(01 章开头的「调用不执行」),此刻上游还一步没跑;async for每次要一个事件,上游才被推动着往前跑一小段——消费节奏决定生产节奏。所以「检查取消标志」能插在每两个事件之间:用户点停止后,最多再跑一个事件就会取消;- 如果把这里的
async for改成[x async for x in ...](列表推导收齐),取消机制立刻失效——因为收齐意味着一口气把生成器跑到底。流式与取消能力,都建立在「逐个消费」上。
三例合观:一条完整的数据流
模型/数据源(原始事件源)
│ async yield
▼
stream 层 ──► stream_events ──► run_pipeline ← 例一(生产者/中转)
(转换) (备料/旁路) (标记取消/落库)
│ async yield ×N 层
▼
Worker async for ← 例三(消费者)
│ 每消费一个:查取消 → 写存储
▼
sse_generator(异步生成器) ← 例二(读存储的生产者)
│ async yield(SSE 帧)
▼
StreamingResponse ──► 浏览器
每一层都用了同一套语言机制(惰性、冻结/恢复、await 让出、break 收尾),区别只是站在管道的哪个位置。
↓下一步:05 章 · 陷阱清单与自检——把容易翻车的十个点过一遍。