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()                        # 无论成败都触发

这个例子的三个知识点:

  1. 「消费上游 + 产出下游」的管道形态:async for item in 上游 + yield item,自己既是消费者又是生产者。事件链路里这种「中转层」往往有好几层,每层只加一点自己的逻辑(这里就是标记取消/错误)——每层都很薄,串起来才是完整管道;
  2. yield item 在 try 块里:下游消费者出异常时,异常会从 yield 处抛进本函数,被 except 接住走兜底——这是 02 章 throw 机制的自动版(消费者不礼貌退出时,生产者能感知);
  3. 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

消费者视角的三个知识点:

  1. await build_event_source(...) 返回的是生成器对象(01 章开头的「调用不执行」),此刻上游还一步没跑;
  2. async for 每次要一个事件,上游才被推动着往前跑一小段——消费节奏决定生产节奏。所以「检查取消标志」能插在每两个事件之间:用户点停止后,最多再跑一个事件就会取消;
  3. 如果把这里的 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 章 · 陷阱清单与自检——把容易翻车的十个点过一遍。

进入 keel 阅读