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

生产环境的 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 章 · 三个真实场景生成器拆解——把这三章的知识对到真实代码上。

进入 keel 阅读