KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
03 · 异步生成器:`async def` + `yield` — keel 龙骨
01 章的生成器解决「惰性生产」,02 章的 send/throw 解决「双向通信」——但它们有个共同限制:函数体里不能 await。模型调用要 await、Redis 读要 await、这些 IO 等待期间函数体必须让出控制权。把 yield 放进 async def,两者就合一了。生产级流式链路里的每一个生成器,都是异步生成器。
01 章的生成器解决「惰性生产」,02 章的 send/throw 解决「双向通信」——但它们有个共同限制:函数体里不能 await。模型调用要 await、Redis 读要 await、这些 IO 等待期间函数体必须让出控制权。把 yield 放进
async def,两者就合一了。生产级流式链路里的每一个生成器,都是异步生成器。
一、先解决一个困惑:async 生成器 vs 协程,别混
两种 async def 长得像,本质完全不同:
async def coroutine(): # 协程(coroutine):没有 yield
return 42
async def agen(): # 异步生成器(async generator):有 yield
yield 1
yield 2
| 协程 | 异步生成器 | |
|---|---|---|
| 调用后得到 | coroutine 对象,必须 await | async generator 对象,必须 async for |
| 怎么拿结果 | x = await coroutine() |
async for item in agen(): ... |
await f() 时 f 是生成器会怎样 |
TypeError:object async_generator can't be used in 'await' expression | —— |
| 执行模型 | 跑到 return/结束 | 跑到下一个 yield 冻结(01 章那套,不变) |
判断口诀:函数体里有 yield → 生成器(要用 async for / async next);没有 yield → 协程(要用 await)。 看到报错 "can't be used in await",第一反应应该是「我是不是把异步生成器当协程 await 了」。
二、异步生成的「异步」体现在哪:yield 之间可以 await
普通生成器的暂停点只有 yield;异步生成器在两个 yield 之间还可以有 await:
import asyncio
async def poll_events():
offset = 0
while True:
records = await fetch_from_redis(offset) # ← IO 等待,期间事件循环去干别的
for r in records:
yield r # ← 产出一个,冻结
offset = r.offset
await asyncio.sleep(0.1) # 没数据时歇一会
它的执行模型 = 01 章的冻结/恢复 + asyncio 的挂起/调度:
async for item in poll_events():
│
├─ 事件循环运行 poll_events 直到 yield r
│ (期间可能 await 过:fetch_from_redis 的网络等待中,事件循环可以跑其他任务)
├─ item = r,循环体执行
└─ 循环体结束时,驱动器调 __anext__() → 从冻结处恢复
两个世界的「暂停」是叠加的:await 是「让出 CPU 给事件循环,等 IO 完成叫我」,yield 是「把值交给消费者,等消费者要下一个叫我」。异步生成器 = 会做 IO 的惰性生产者。
2.1 手动驱动:__anext__
01 章用 next() 手动驱动同步生成器;异步版要 await gen.__anext__():
g = poll_events()
first = await g.__anext__() # 手动要一个(不能用 next()!)
async for 的翻译版本(对照 01 章 3 节):
# async for item in agen(): 大致等价于
it = agen().__aiter__()
while True:
try:
item = await it.__anext__()
except StopAsyncIteration: # 注意:不是 StopIteration!
break
...
易错点:异步迭代器的结束信号是 StopAsyncIteration,不是 StopIteration。框架代码里看到捕获哪一个,就知道它消费的是同步还是异步生成器。
三、和「返回列表的协程」对比:为什么流式必须用异步生成器
把 00 章的对话例子搬进异步世界,对比就清楚了:
# 方案 A:协程返回列表
async def chat_all():
pieces = []
async for token in model_stream():
pieces.append(token)
return pieces # 调用方:response = await chat_all() → 30 秒后一次性拿到
# 方案 B:异步生成器
async def chat_stream():
async for token in model_stream():
yield token # 调用方:async for piece in chat_stream() → 每个片段即时可用
方案 B 的每个 token 在产出的瞬间就到了消费者手里,消费者(SSE handler)立刻能 yield f"data: ..." 推给浏览器。链路形态:
model loop(异步产出原始 token)
→ transform_events(异步生成器:转换事件) ← yield
→ enrich(异步生成器:补 ID、落库) ← yield
→ observe(异步生成器:旁路观察) ← yield
→ sse_handler(异步生成器:读队列 + SSE 帧) ← yield
→ StreamingResponse(消费,推给浏览器)
五层异步生成器首尾相接,token 从模型流到浏览器。任何一层用了「攒列表再返回」,整条链就退化为「等最慢的一层攒完」——这就是 00 章说 yield 是「流式系统的泵」的含义。
四、异步生成器的资源清理:async with / aclose()
02 章的 close/GeneratorExit 在异步版同样存在,配套换成 async 版本:
async def sse_events(request):
try:
while True:
if await request.is_disconnected(): # 每轮检查客户端
break
...
yield frame
finally:
# 客户端断开 / 生成器被 close 时执行清理
await release_resources()
await g.aclose()在冻结的 yield 处抛GeneratorExit;- 生成器里可以有
async with/await形式的清理(finally 里 await 是合法的,同步生成器做不到); - 消费方中途退出 async for(break/return/异常),生成器被 GC 时同样会 aclose——但显式管理更可靠。
生产环境的 SSE 生成器(04 章逐行走)就是靠 request.is_disconnected() 主动 break 来结束自己——因为 StreamingResponse 框架会在断开时停止拉取,生成器下一轮检查到断开就退出。「停止消费 → 停止生产」在异步世界依然成立。
五、类型注解速查(写代码用)
from collections.abc import AsyncGenerator, AsyncIterator
from typing import Any
async def agen() -> AsyncGenerator[int, None]:
# ↑yield 的类型 ↑send 的类型(通常 None)
...
# AsyncIterator 是它的「只读」版本,消费方参数常用:
async def consume(it: AsyncIterator[Any]) -> None:
async for x in it:
...
↓下一步:04 章 · 三个真实场景生成器拆解——把这三章的知识对到真实代码上。