[LangChain]聊天模型-深入 LangChain 流式处理与异步架构的底层逻辑
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. 执行流程与输出
事件循环会按以下逻辑调度:
- 执行
boil_water()→ 打印“开始烧水...” → 遇到await asyncio.sleep(5),暂停,转去执行send_message()。- 执行
send_message()→ 打印“开始发短信...” → 遇到await asyncio.sleep(2),暂停。- 2秒后,
send_message()的sleep完成 → 事件循环恢复它 → 打印“短信发送成功!” → 任务结束。- 再过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)['在星空下许下心愿的承诺']|['你的笑容是我永恒的向往']|['默默相依,岁月无声流淌']|['爱在时光里绽放,愈发芬芳']|['只愿与你携手,共度每个晌午']|现在的流程变成了:
- Model: 吐出一个字 "今"
- Parser: 变成字符串 "今"
- Splitter: 存入 buffer ("今") -> 没句号 -> 等待
- ... Model 继续吐 ...
- 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()。 - 内部机制:
- LangChain 调用 OpenAI SDK。
- SDK 发起 HTTP 请求,并在请求参数中带上
stream=True。 - OpenAI 服务器建立 SSE 连接,开始源源不断地返回数据块。
- 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等字段)- 处理过程:
- 提取
delta中的内容- 判断是否有
function_call或tool_calls(工具调用)- 根据角色(role)和内容(content),实例化不同的 Chunk 类(如
ChatGenerationChunk)
- 输出: 一个标准的
AIMessageChunk对象,你的代码最终拿到的就是这个
总结 :
- 本质:基于 HTTP 的长连接。服务器不关闭连接,而是源源不断地发送数据块
- 格式:非常简单的文本格式,以
data:开头,以\n\n结束 - 对比 WebSocket:SSE 是单向的(服务器 -> 客户端),非常适合 LLM 生成场景;WebSocket 是双向的,更适合聊天室
深度解析:LangChain 的 astream_events 方法之所以强大,是因为它不仅流式传输“最终文本”,还流式传输“中间状态”(比如工具调用的开始、结束、LLM 的思考过程)。这让前端可以做出非常炫酷的“打字机 + 状态指示器”效果
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐





所有评论(0)