1.核心痛点:为什么需要流式传输?

在 AI 应用开发中,最大的痛点就是“延迟”。当 LLM 生成几百字的内容时,如果让用户盯着空白屏幕等几十秒,体验是灾难性的

传统模式 : 请求 -> 等待 -> 完整返回

流式传输 : 请求 -> 逐字/逐块返回

流式传输不仅是为了“快”,更是为了掩盖延迟。流式传输让系统看起来一直在思考和工作,极大提升交互体验

2. 同步 vs 异步

2.1 stream() 同步

from langchain_openai import ChatOpenAI
# 定义大模型
model = ChatOpenAI(model="gpt-4o-mini")
# 流式输出
chunks = []
for chunk in model.stream("讲一个50字的笑话"):
    chunks.append(chunk)
    print(chunk.content, end="|", flush=True)
|有|一天|,小|明|问|老师|:“|如果|我|把|你的|书|借|给|别人|,|算|不|算|偷|呢|?”|老师|回答|:“|当然|算|!”|小|明|一|脸|困|惑|:“|那|我|借|你的|知识|呢|?”|老师|笑|着|说|:“|那|叫|共享|智慧|,|绝|对|不|算|偷|!”|||
  • 机制:使用 for 循环遍历 model.stream()
  • 缺点:代码虽然简单,但它是阻塞的。在处理当前 token 时,线程无法做其他事。如果是 Web 服务,一个请求就会占死一个线程。

2.2 astream()异步

异步方式:需要用到 asyncio,协程,事件循环

① 协程(Coroutine)

是 python 中实现单线程内并发的核心载体,本质是可在执行过程中“暂停”和“恢复”的特殊函数

定义与语法
  • async def 声明协程函数(区别于普通函数 def)。
  • await 关键字标记“暂停点”:当协程执行到 await 时,会主动暂停自身,把控制权交还给事件循环,等待异步操作(如网络请求、文件读写、定时等待)完成后,再恢复执行。
与进程、线程的区别

  • 进程:多核 CPU 并行,每个进程有独立内存空间,切换开销大(需操作系统调度)。
  • 线程:单核 CPU 并发(时间片轮转),共享进程内存,但线程切换仍有内核级开销。
  • 协程:在单个线程内实现“伪并行”,由用户代码(而非操作系统)控制暂停/恢复,切换开销极小(本质是函数栈的保存与恢复)
代码示例 :
import asyncio

# 异步函数:模拟烧水
async def boil_water():
    print("开始烧水...")
    # 异步休眠5秒,协程主动让出控制权,事件循环可以调度其他任务
    await asyncio.sleep(5)
    print("水开了!")

# 异步函数:模拟发送短信
async def send_message():
    print("开始发短信...")
    # 异步休眠2秒,协程暂停,释放事件循环
    await asyncio.sleep(2)
    print("短信发送成功!")

# 主异步入口函数
async def main():
    # 创建任务,将协程交给事件循环调度,实现并发执行
    task1 = asyncio.create_task(boil_water())
    task2 = asyncio.create_task(send_message())
    # 等待烧水任务执行完成
    await task1
    # 等待发短信任务执行完成
    await task2

# 程序入口
if __name__ == "__main__":
    # asyncio.run() 创建事件循环,运行main协程,结束后自动关闭循环
    asyncio.run(main())

② 事件循环(Event Loop)

事件循环是 asyncio 的核心引擎,负责调度所有协程的执行顺序,让单线程能同时处理多个异步任务

1. 工作原理

事件循环维护一个任务队列,不断循环做两件事:

  • 检查任务状态:如果某个协程处于“等待异步操作完成”(如 await asyncio.sleep()),就把它暂时挂起,去执行其他“就绪”的协程。
  • 恢复已完成等待的任务:当异步操作完成(如 sleep 时间到、网络请求返回),事件循环会把对应的协程重新放回队列,继续执行后续代码。
2. 启动与运行

通过 asyncio.run(主协程()) 启动事件循环。run() 会自动创建事件循环,运行传入的主协程,并在主协程结束后关闭循环。

3. 执行流程与输出

事件循环会按以下逻辑调度:

  1. 执行 boil_water() → 打印“开始烧水...” → 遇到 await asyncio.sleep(5)暂停,转去执行 send_message()
  2. 执行 send_message() → 打印“开始发短信...” → 遇到 await asyncio.sleep(2)暂停
  3. 2秒后,send_message()sleep 完成 → 事件循环恢复它 → 打印“短信发送成功!” → 任务结束。
  4. 再过3秒(累计5秒),boil_water()sleep 完成 → 事件循环恢复它 → 打印“水开了!” → 任务结束。

最终输出(总耗时约5秒,而非 5+2=7秒,因为两个 sleep并发等待):

1开始烧水...
2开始发短信...
3短信发送成功!
4水开了!

    使用 :

    import asyncio
    
    from langchain_openai import ChatOpenAI
    
    model = ChatOpenAI(model = "gpt-4o-mini")
    
    # 异步调用
    async def async_stream():
        print("异步调用")
        async for chunk in model.astream("讲一个50字的笑话"):
            print(chunk)
            
    asyncio.run(async_stream())

    异步调用
    |有|一天|,小|明|问|老师|:“|为什么|我们的|书|总|是|那么|厚|?”|老师|回答|:“|因为|知识|就是|力量|!”|小|明|陷|入|沉|思|:“|那|我|是不是|可以|锻|炼|身体|,|吃|更多|的|书|?”|老师|无|奈|一|笑|:“|可|别|,|书|可|不能|当|饭|吃|!”|||
    进程已结束,退出代码为 0

    总的来说:

    • 协程是“可暂停的函数”,是异步任务的载体
    • 事件循环是“调度中枢”,负责决定哪个协程该执行、哪个该等待
    • asyncio是封装了协程和事件循环的工具库,让异步 I/O 开发更高效、更易读

    2.3 使用 StrOutputParser 解析流式输出

    大模型在流式输出时,返回的通常不是纯文本字符串,而是包含元数据的对象(如 AIMessageChunk)。我们需要把这些对象“清洗”成人类可读的字符串

    • 工具StrOutputParser
    • 作用:从模型返回的复杂对象中提取出纯文本内容
    • 代码逻辑
    # 1. 定义模型
    model = ChatOpenAI(model="gpt-4o-mini")
    # 2. 定义解析器
    parser = StrOutputParser()
    # 3. 组装链条 (Chain)
    chain = model | parser
    # 4. 流式调用
    for chunk in chain.stream("写一段关于爱情的歌词..."):
        print(chunk, end="|", flush=True)
    |(|Verse| |1|)|  
    |在|星|空|下|许|下|承|诺|,|  
    |月|光|洒|在|你|我的|肩|头|。|  
    |心|跳|声|中|交|织|着|温|柔|,|  
    |每|个|瞬|间|都|如|梦|似|幻|的|柔|。
    
    |(|Ch|orus|)|  
    |爱|是|花|海|中的|芬|芳|,|  
    |在|风|中|飘|荡|,|愈|加|迷|醉|。|  
    |与你|相|拥|,|忘|却|时|光|,|  
    |这|份|心|意|永|远|不|变|,|无限|蔓|延|。
    
    |(|Verse| |2|)|  
    |手|中的|温|度|藏|着|未来|,|  
    |眼|神|交|汇|似|海|般|深|邃|。|  
    |每|一次|触|碰|都是|期待|,|  
    |将|这一|刻|镌|刻|心|间|不|再|离|开|。
    
    |(|Ch|orus|)|  
    |爱|是|花|海|中的|芬|芳|,|  
    |在|风|中|飘|荡|,|愈|加|迷|醉|。|  
    |与你|相|拥|,|忘|却|时|光|,|  
    |这|份|心|意|永|远|不|变|,|无限|蔓|延|。
    
    |(|Bridge|)|  
    |即|使|风|雨|再|大|,|  
    |我|也|与你|并|肩|,|  
    |心|中的|光|芒|,|  
    |永|远|指|引|前|行|的|方向|。
    
    |(|Ch|orus|)|  
    |爱|是|花|海|中的|芬|芳|,|  
    |在|风|中|飘|荡|,|愈|加|迷|醉|。|  
    |与你|相|拥|,|忘|却|时|光|,|  
    |这|份|心|意|永|远|不|变|,|直到|天|荒|地|老|。|||
    进程已结束,退出代码为 0

    2.4 自定义流式输出解析器(核心难点)

    ① 场景需求

    我们希望模型一边生成内容,我们一边处理,但不是按“字”处理,而是凑齐一个完整的句子(以句号结尾)后再输出或处理。

    ② 实现原理:生成器函数

    为了实现这个功能,代码定义了一个名为 split_into_list 的函数。注意它的类型签名:
    Iterator[Input] -> Iterator[List[str]]

    这意味着它是一个生成器,它接收上游传来的数据流,处理后,再向下游抛出新的数据流。

    ③ 代码逻辑深度解析

    from langchain_openai import ChatOpenAI
    from langchain_core.output_parsers import StrOutputParser
    from typing import Iterator, List
    
    # 1. 定义模型
    model = ChatOpenAI(model="gpt-4o-mini")
    # 2. 定义解析器
    parser = StrOutputParser()
    
    def split_into_list(input: Iterator[str]) -> Iterator[List[str]]:
        buffer = ""  # 缓冲区:用来暂存还没凑成句子的碎片
    
        # 1. 遍历上游传来的每一个小碎片 (chunk)
        for chunk in input:
            buffer += chunk  # 把碎片拼接到缓冲区
    
            # 2. 检查缓冲区内是否有句号 "."
            while "。" in buffer:
                # 找到第一个句号的位置
                stop_index = buffer.index("。")
    
                # 【关键动作 yield】:
                # 截取句号前的内容作为一个完整的句子,抛给下游
                # strip() 去除首尾空格
                yield [buffer[:stop_index].strip()]
    
                # 更新缓冲区:把已经抛出的句子切掉,保留剩下的部分
                buffer = buffer[stop_index + 1:]
    
        # 3. 循环结束后,如果缓冲区里还有残留文字(没凑够一句),也抛出去
        if buffer:
            yield [buffer.strip()]
    # 3. 组装
    chain = model | parser | split_into_list
    
    # 4. 流式调用
    for chunk in chain.stream("写一段关于爱情的歌词,需要 5句话,每句话需要句号分割"):
        print(chunk, end="|", flush=True)
    ['在星空下许下心愿的承诺']|['你的笑容是我永恒的向往']|['默默相依,岁月无声流淌']|['爱在时光里绽放,愈发芬芳']|['只愿与你携手,共度每个晌午']|

    现在的流程变成了:

    1. Model: 吐出一个字 "今"
    2. Parser: 变成字符串 "今"
    3. Splitter: 存入 buffer ("今") -> 没句号 -> 等待
    4. ... Model 继续吐 ...
    5. Splitter: buffer 变成 "今天天气不错." -> 发现句号 -> Yield ["今天天气不错"] -> 清空 buffer

    3.SSE 协议介绍(server-sent events)

    1. 这是所有流式传输的基石

    ①什么是 SSE?

    • 它是一种基于 HTTP 协议的单向通信技术。
    • 特点:服务器可以持续向客户端推送数据,而不需要客户端反复请求(轮询)。
    • 为什么用它? LLM 生成文本是逐字生成的,SSE 允许服务器每生成一个字(Token)就立刻发给客户端,从而实现“打字机”效果。

    ②协议细节

    • Content-Type: text/event-stream
    • Connection: keep-alive (保持连接不断开)
    • 数据格式: 每次发送的数据块以 data: 开头,块与块之间用 \n\n 分隔。

    2. LangChain 的流式传输流程分析

    • 入口方法: stream()astream()
    • 内部机制:
    1. LangChain 调用 OpenAI SDK。
    2. SDK 发起 HTTP 请求,并在请求参数中带上 stream=True
    3. OpenAI 服务器建立 SSE 连接,开始源源不断地返回数据块。
    4. LangChain 接收到这些原始数据块后,将其转换为统一的 AIMessageChunk 对象。

    3. 源码级深度解析:OpenAI SDK 是如何工作的?

    A. 客户端封装

    LangChain 使用了一个叫 _SyncHttpxClientWrapper 的类来包装 OpenAI 的 HTTP 客户端。这确保了网络请求的稳定性和超时控制

    B. 关键转换逻辑:_convert_chunk_to_generation_chunk

    这是最核心的代码段。OpenAI 返回的是原始的 JSON 字典(Raw Dict),LangChain 需要把它变成自己能懂的对象

    • 输入: OpenAI 返回的原始 chunk(包含 choices, delta, content 等字段)
    • 处理过程:
    1. 提取 delta 中的内容
    2. 判断是否有 function_calltool_calls(工具调用)
    3. 根据角色(role)和内容(content),实例化不同的 Chunk 类(如 ChatGenerationChunk
    • 输出: 一个标准的 AIMessageChunk 对象,你的代码最终拿到的就是这个

    总结 :

    • 本质:基于 HTTP 的长连接。服务器不关闭连接,而是源源不断地发送数据块
    • 格式:非常简单的文本格式,以 data: 开头,以 \n\n 结束
    • 对比 WebSocket:SSE 是单向的(服务器 -> 客户端),非常适合 LLM 生成场景;WebSocket 是双向的,更适合聊天室

    深度解析:LangChain 的 astream_events 方法之所以强大,是因为它不仅流式传输“最终文本”,还流式传输“中间状态”(比如工具调用的开始、结束、LLM 的思考过程)。这让前端可以做出非常炫酷的“打字机 + 状态指示器”效果

    Logo

    openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

    更多推荐