news 2026/9/30 6:39:45

LangChain 流式输出完全指南:从 AIMessageChunk 到 stream_events 的 Token 级流式架构与实战

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
LangChain 流式输出完全指南:从 AIMessageChunk 到 stream_events 的 Token 级流式架构与实战
  • 人工智能
  • 大模型
  • AI Agent
  • Agent 框架
  • RAG

【免费下载链接】langchain

The agent engineering platform.

项目地址:https://gitcode.com/GitHub_Trending/la/langchain
点击查看免费下载

Streaming(流式输出)是 LangChain 将模型输出按 Token 逐块交付的核心机制:应用不再像invoke()那样阻塞等待完整响应,而是通过stream()/astream()持续接收增量结果,实现 Web 界面、控制台等场景下的实时反馈。本文以 openwiki/streaming.md 为骨架,结合langchain_core源码,系统讲解同步/异步流式协议、AIMessageChunk的增量合并、回调集成、链式流式传播、stream_events(version="v3")的事件流与反压机制,以及流式与invoke()的内存/延迟权衡。读完本文,你将掌握在 LangChain 应用中搭建 Token 级实时输出、结构化流式解析与 Agent 执行可视化的完整方案。

概述:什么是流式输出

流式(Streaming)是指 LangChain 将模型输出增量交付(token by token),而不是等待整个响应生成完毕后再一次性返回。它带来两方面的价值:

  • 实时反馈:Web UI、控制台等面向用户的应用可以逐字显示输出,不必在模型延迟期间白屏等待;
  • 响应式应用基础:为不阻塞在模型延迟上的交互式应用提供底层支撑。

其核心机制是:应用调用stream()或astream(),获得一个由部分输出组成的序列(每个部分是一个携带增量内容的AIMessageChunk);回调通过on_llm_new_token事件拦截这些分块,从而在不必收集完整响应的前提下,对每个 Token 进行观察、日志记录或即时响应。流式同样贯穿于由 Prompt、模型、输出解析器等 Runnable 组成的链条(chain)中,只要链的每个组件都支持流式,组合出来的链就自动支持流式。

同步流式:stream()

位置:libs/core/langchain_core/language_models/chat_models.py

BaseChatModel.stream()是同步流式的主要入口,它按底层模型产出的顺序 yield 出AIMessageChunk对象,每个分块携带增量内容——可能是一个 Token、一段 JSON 片段或一次结构化块更新。

控制流

stream()的完整执行流程如下:

  1. 检查模型是否支持流式:_should_stream()决定本次调用是否走流式路径(chat_models.py#L549-L585),判断依据包括:

    • _stream()是否在模型子类上真正实现(而非继承自BaseChatModel基类);
    • 是否通过disable_streaming、stream=False或streaming=False显式禁用了流式;
    • 是否显式传入stream=True;
    • 模型上是否挂接了流式回调处理器(handler 是否为_StreamingCallbackHandler实例)。

    如果模型不支持流式,stream()会回退到invoke(),并将完整结果强制转换为AIMessageChunk后一次性 yield 出来。

  2. 初始化回调:根据传入的RunnableConfig配置一个CallbackManager,绑定 callbacks、tags 与 metadata,用于链路追踪与可观测性。

  3. 触发on_chat_model_start:回调生命周期从on_chat_model_start开始,宣告 LLM 调用即将开始,该事件在第一个 Token 产出之前触发。

  4. 获取限流许可:如果模型挂接了 rate limiter,stream()会先rate_limiter.acquire(blocking=True)获取许可,在限流允许前阻塞等待(chat_models.py#L785-L786)。

  5. 迭代模型分块:对底层_stream()产出的每个ChatGenerationChunk:

    • 若分块消息 ID 为空,则将其设置为以LC_ID_PREFIX前缀的唯一 run ID,保证可追踪性;
    • 通过_gen_info_and_msg_metadata()计算并附加响应元数据(模型提供商、延迟、Token 用量等);
    • 若模型output_version == "v1"(content-block 结构化格式),则将内容转换为内容块并按type索引编号;
    • 触发on_llm_new_token:携带该分块的内容与完整分块对象,让回调可以观察或缓冲每个 Token;
    • 将分块消息转换为AIMessageChunk并立即 yield 给调用方;
    • 同时把分块累积到内存列表,供后续聚合使用。
  6. yield 最终的 "last" 分块:模型结束后,若没有显式设置chunk_position="last",则 yield 一个内容为空、chunk_position="last"的末尾分块(chat_models.py#L816-L832)。该信号通知解析器与消费者流已结束,同时触发tool_call_chunks终结为完整的tool_calls。

  7. 回调生命周期收尾:成功时触发on_llm_end,携带合并了所有分块的ChatGeneration;若发生异常则触发on_llm_error(携带部分累积结果),让回调在异常重新抛出前观察到失败现场(chat_models.py#L833-L855)。

回退行为

如果模型未实现流式(通过_should_stream(async_api=False)判断),stream()会委托给invoke()并 yield 单个强制转换为AIMessageChunk的结果。这一设计保证了所有模型都提供一致的流式接口,即便底层只支持非流式的invoke()。_should_stream的实现(chat_models.py#L559-L585)还兼顾了异步场景:异步未实现但同步已实现时,异步路径可以回退到同步流式。

异步流式:astream()

位置:libs/core/langchain_core/language_models/chat_models.py

BaseChatModel.astream()是stream()的异步变体,逻辑与同步版本一一对应,但使用async/await与AsyncCallbackManager。关键差异如下:

  • 回调事件使用await:await run_manager.on_llm_new_token(...)、await run_manager.on_llm_end(...);
  • 通过async for chunk in self._astream(...)迭代底层异步流式实现;
  • 限流许可通过await self.rate_limiter.aacquire(blocking=True)获取(chat_models.py#L916-L917)。

异步流式协议与同步完全一致:检查_should_stream(async_api=True)→ 初始化回调 → 分块到达即 yield → 每 Token 触发回调 → 在 "last" 信号上终结 tool call chunks。

Async/Await 使用模式

使用astream()的应用应在循环中消费异步迭代器:

async for chunk in model.astream(messages): # 立即处理分块 print(chunk.content, end="", flush=True) # 或者收集分块供后续处理 chunks = [] async for chunk in model.astream(messages): chunks.append(chunk) final_message = sum(chunks) # 通过 + 运算符合并

AIMessageChunk:增量内容

位置:libs/core/langchain_core/messages/ai.py

AIMessageChunk是流式过程中 yield 的消息类型。与完整的AIMessage不同,它表示对会话消息的部分、增量更新,并支持通过+运算符合并。

结构

  • content:字符串或内容块列表(当output_version="v1"时)。流式期间每个分块只包含该步的新 Token 或增量;内容跨分块累积——文本分块拼接,JSON 分块可能追加部分对象或数组。
  • tool_call_chunks:ToolCallChunk对象列表(正在流式传输的不完整工具调用)。模型产出调用 ID、函数名与参数 JSON 时被逐步更新;参数通过parse_partial_json()增量累积与解析。
  • chunk_position:可选哨兵值;设为"last"表示流中的最后一个分块,触发工具调用与推理块的终结。该分块被聚合时,tool_call_chunks会经init_tool_calls()校验器解析为完整的tool_calls与invalid_tool_calls(ai.py#L508-L601)。
  • response_metadata:模型相关元数据(延迟、model_provider、用量计数器、finish reason 等),由流式处理器附加;元数据跨分块合并,其中 usage 计数做求和。

合并与聚合

流式分块通过+运算符累积,其底层实现为add_ai_message_chunks()(ai.py#L652-L700),合并逻辑如下:

  1. 合并内容:文本内容拼接;结构化内容块按块类型合并。
  2. 拼接工具调用参数:tool_call_chunks的参数被逐步追加,支持增量 JSON 解析。
  3. 合并元数据:响应元数据与用量计数合并;重复键时后值覆盖(usage 例外,做求和)。
  4. 保留 chunk_position:合并中的任一分块若带有chunk_position="last",结果即标记为 "last",触发工具调用终结。
  5. 择优选择 ID:分块 ID 按优先级排序——提供商分配的(非LC_*前缀)>LC_run_*>lc_*自动 ID。

当分块被合并或收到 "last" 信号时,会重建一个tool_calls已终结(而非 chunk)的完整AIMessage:

# 累积分块 chunks = [] async for chunk in model.astream(messages): chunks.append(chunk) # 将所有分块合并为一条消息 final_message = chunks[0] for chunk in chunks[1:]: final_message = final_message + chunk # tool_calls 现在是完整的 ToolCall 对象,而非 ToolCallChunk for tool_call in final_message.tool_calls: print(tool_call["name"], tool_call["args"])

回调集成:on_llm_new_token

位置:libs/core/langchain_core/callbacks/base.py

on_llm_new_token回调在流式过程中为每个 Token 或分块触发一次,是实时观察与日志记录的核心钩子。它定义在LLMManagerMixin中,同时适用于聊天模型与传统文本补全 LLM。

签名

def on_llm_new_token( self, token: str | list[str | dict[str, Any]], *, chunk: GenerationChunk | ChatGenerationChunk | None = None, run_id: UUID, parent_run_id: UUID | None = None, tags: list[str] | None = None, **kwargs: Any, ) -> Any:

参数说明:

  • token:字符串 Token 或内容块列表(output_version="v1"时)。文本流式时是单个词或子词;结构化输出时是内容块 dict 列表,包含type、text、reasoning、tool_call_chunk等字段。
  • chunk:完整的ChatGenerationChunk,携带元数据、消息 ID、响应元数据与tool_call_chunks。这让回调可以检查完整的分块结构,而不仅是 Token 本身。
  • run_id:本次流式运行的唯一标识,用于追踪以及与父操作的关联。
  • parent_run_id:调用该模型的父 run(链或 Agent)的 ID。
  • tags:来自调用上下文的可继承标签,可用于过滤或路由回调。

示例:流式输出到 stdout

from langchain_core.callbacks import StreamingStdOutCallbackHandler callback = StreamingStdOutCallbackHandler() # 回调通过 RunnableConfig 传入 for chunk in model.stream( messages, config=RunnableConfig(callbacks=[callback]) ): pass # 回调将每个 token 打印到 stdout

StreamingStdOutCallbackHandler(libs/core/langchain_core/callbacks/streaming_stdout.py)实现了on_llm_new_token,将 Token 立即写入sys.stdout并flush(),使流式输出无需缓冲即可实时可见。注意其文档字符串中的警告:仅适用于支持流式的 LLM。

自定义流式回调

通过继承BaseCallbackHandler创建自定义回调:

from langchain_core.callbacks import BaseCallbackHandler class MyStreamingCallback(BaseCallbackHandler): def on_llm_new_token(self, token: str, **kwargs: Any) -> None: # 将 token 发送到 WebSocket、写入数据库等 websocket.send_json({"token": token})

流式穿过链(Streaming Through Chains)

流式贯穿由 Runnable(提示词、模型、解析器)组成的链。流式协议在每一层通过Runnable的stream()与transform()方法实现。

默认行为

位置:libs/core/langchain_core/runnables/base.py

默认情况下,Runnable.stream()只 yield 一次invoke()的完整输出(base.py#L1194-L1213);astream()同理回退到ainvoke()。支持流式的子类必须覆写stream()或transform()来产出分块。transform()是核心流式接口:它接收输入迭代器并产出输出迭代器,从而支持有状态的逐块变换。

RunnableSequence 中的流式

位置:libs/core/langchain_core/runnables/base.py

RunnableSequence(用|运算符创建的链)在满足以下条件时自动支持流式:

  1. 所有上游组件实现了transform():transform()将流式输入映射为流式输出,实现端到端无缓冲流式;
  2. 最后一个组件能产出分块:输出解析器与模型实现transform()以产出部分结果。

若链中任一组件未实现transform(),流式将在该组件处阻断(形成缓冲点)。多个阻断组件会形成多个缓冲点,但只要末组件支持流式,最终输出仍会从该组件开始流式。

重要:RunnableLambda默认不实现transform(),因此是阻断组件。若需要自定义逻辑参与流式,应继承Runnable并覆写transform()。

流式示例:Model → Parser

from langchain_openai import ChatOpenAI from langchain_core.output_parsers import StrOutputParser model = ChatOpenAI() parser = StrOutputParser() chain = model | parser # stream 按 token 到达的顺序产出解析器输出 for chunk in chain.stream("What is 2+2?"): print(chunk, end="", flush=True)

当model.stream()产出分块时,解析器的transform()(继承自BaseTransformOutputParser)消费每个分块并产出其变换结果。文本解析器(如StrOutputParser)从AIMessageChunk中抽取文本并直接 yield 字符串;JSON 解析器则通过parse_partial_json()在分块可解析时产出部分 JSON 对象。

流式机制:_transform_stream_with_config

位置:libs/core/langchain_core/runnables/base.py

_transform_stream_with_config()帮助函数负责管理带回调的流式转换,流程如下:

  1. Tee 输入迭代器:先窥视第一个元素用于追踪,而不消费它;
  2. 触发on_chain_start:在转换器开始前触发,宣告链式流式操作的开始;
  3. 调用转换器函数:传入剩余输入迭代器与子回调;
  4. 立即 yield 分块:转换器产出即交付,保证流式响应;
  5. 累积输出:为on_chain_end回调累积输出(若支持则通过+合并分块);
  6. 触发on_chain_end或on_chain_error:结束时提供合并后的最终输出或异常上下文。

这一机制保证每个分块都会触发链式流式回调,并让父 run manager 知道链的流式何时结束。

通过stream_events流式:ChatModelStream

位置:libs/core/langchain_core/language_models/chat_model_stream.py

对于需要细粒度事件的高级场景,BaseChatModel.stream_events(version="v3")(chat_models.py#L1286-L1358)返回一个ChatModelStream对象,它暴露类型化投影属性(.text、.tool_calls、.usage、.reasoning、.output),随协议事件到达而持续累积。异步对应物为astream_events(version="v3")返回的AsyncChatModelStream。

需要说明的是,version="v3"目前处于 beta 阶段:协议形态、返回类型与 API 面积在未来版本中可能变化,调用时会发出LangChainBetaWarning。此外,v3 的输出ChatModelStream.output.content始终是 v1 内容块列表(text / reasoning / tool_call / image / …),与模型自身的output_version属性无关;若在同一流水线中混用 v3 与传统stream()/invoke()并需要一致的输出形态,应在模型上设置output_version="v1"。

结构化事件流

与 yield Token 的stream()不同,stream_events()yield 的是协议事件——表示模型状态变化的结构化对象:

  • text-delta:增量文本生成;
  • reasoning-delta:增量推理/思考内容(当模型支持时);
  • tool_call_chunk:带累积参数的部分工具调用;
  • usage:Token 用量更新(输入、输出、缓存等)。

基于拉取的反压(Pull-Based Backpressure)

ChatModelStream及其投影(.text、.tool_calls等)通过SyncProjection与AsyncProjection类实现基于拉取的反压。当消费者从投影读取并追平缓冲时:

  1. 投影调用_request_more()从生产者(模型/图)拉取额外事件;
  2. 生产者恢复并生成下一批事件;
  3. 投影缓冲事件并 yield 给消费者。

该反压机制防止无界内存增长:生产者只在消费者请求时生成事件。与回调驱动的流式(模型产出多快就交付多快)不同,基于拉取的流式允许消费者掌控节奏。

示例:带反压的消费

# 按需拉取事件;若无人拉取,生产者等待 for event in model.stream_events(messages, version="v3"): if should_stop_early(): break # 生产者停止;不会累积缓冲事件 process_event(event) # 或者消费特定投影以获得类型安全 stream = model.stream_events(messages, version="v3") for text_delta in stream.text: # 仅文本事件 print(text_delta)

内存与延迟权衡:stream()vsinvoke()

invoke()

  • 延迟:等待完整模型响应返回后才交付,延迟等于完整的模型生成时间。
  • 内存:无需中间存储,仅持有最终消息。
  • 响应性:阻塞调用线程/协程直至完成,用户在整个响应生成完之前看不到任何输出。
  • 适用场景:批处理;需要先拿到完整响应才能继续下一步的情形。

stream()

  • 延迟:第一个 Token 可用即 yield,最大限度缩短首 Token 时间(TTFT)。
  • 内存:若调用方收集分块,则需要缓冲累积的分块;但由于分块即时 yield,调用方可以处理一个丢弃一个,无需持有整个响应。
  • 响应性:非阻塞,支持渐进式显示,用户实时看到输出。
  • 适用场景:Web UI、控制台应用、依赖实时反馈改善体验的用户交互。

流式并不增加延迟

实践中,流式相比invoke()不会增加显著延迟——模型产出 Token 的速率相同。区别仅在于Token 交付给调用方的时间点。对交互式应用而言,流式交付更优:用户看到的是实时出现的输出,而不是在完整响应就绪前面对一片空白。

反压与内存影响

使用stream()流式时:

  • 回调驱动交付:Token 随模型产出尽快 yield。若调用方消费缓慢,Token 会在stream()内部的累积列表中堆积,直到循环结束或下一次 yield;
  • 无无界增长:分块累积器仅用于最终的on_llm_end回调;分块在累积之前就已 yield。因此内存开销与响应大小成正比,而非与模型速度相关;
  • 消费者节奏:慢消费者(如写磁盘)不会产生反压,只是按 yield 的顺序逐个处理 Token。

使用stream_events()(v3)时:

  • 基于拉取的反压:生产者(模型/图)只在消费者通过投影迭代器请求时才生成事件,天然将生产者节奏与消费者对齐;
  • 有界缓冲:投影只缓冲到消费者读取为止。慢消费者会自然拖慢生产者,防止无界内存增长;
  • 多独立消费者:不同投影(如.text与.tool_calls)上的多个for循环可以从缓冲中回放全部事件,支持多种消费模式而无需重新运行模型。

Agent 执行中的流式

Agent 可以通过create_agent()返回的 Agent 图上的stream_events(version="v3")进行流式执行,从而实时观察:

  • 工具调用:通过.tool_calls投影,追踪哪些工具被调用以及参数随流式累积的过程;
  • 工具输出:工具执行的增量输出,包括产出输出增量的工具所流的增量;
  • 消息:逐步演进的完整对话历史;
  • 子图:当 Agent 通过工具调用内部子 Agent 时,子图句柄暴露各自的投影,实现嵌套可见性。

流式模式(langgraph):

  • "updates":yield 节点更新——哪个节点运行了、产出了什么状态;
  • "values":每个节点完成后 yield 完整状态快照;
  • "messages":仅 yield 消息更新;
  • "custom":由工具或中间件通过emit()或runtime.emit_output_delta()触发的用户自定义流事件。

Agent 流式使循环执行与工具交互过程实时可见,而无需阻塞等待整个 Agent 运行结束。

流式最佳实践

  1. 立即冲刷输出:在 Web 或终端展示流式输出时,每个分块后都要 flush 缓冲,确保即时可见(end=""配合flush=True)。
  2. 谨慎处理部分 JSON:JSON 解析器应使用parse_partial_json()在 Token 到达时提取可解析的完整结构,而不是等待完整响应。
  3. 合并分块以获取最终结果:如果需要完整响应,收集分块并用+合并:
    chunks = [chunk for chunk in model.stream(messages)] final = chunks[0] for chunk in chunks[1:]: final = final + chunk
  4. 用回调处理副作用:将日志、指标与 webhook 交给on_llm_new_token实现,而不是在循环里逐个处理 yield 的分块。回调将应用逻辑与流式关注点解耦。
  5. 尊重反压:使用stream_events()时,让消费者为生产者定速,不要人为加速事件生成。
  6. 选择性关闭流式:对长耗时操作或需要可预测延迟的场景,改用invoke(),或传stream=False覆写默认行为。_streaming_disabled()(chat_models.py#L525-L547)提供了完整的关闭语义:disable_streaming=True、disable_streaming="tool_calling"(且传了 tools)、stream=<falsy>以及实例级streaming=False都会硬性禁用流式。
  7. 同步与异步路径都要测试:stream()与astream()的行为可能因模型实现与回调执行器不同而有差异,务必针对自己的场景两条路径都测。

小结

流式是 LangChain 构建响应式 LLM 应用的核心能力:stream()/astream()提供了 Token 级的增量交付,AIMessageChunk通过+合并实现增量内容的无损聚合,on_llm_new_token回调将实时观察与业务逻辑解耦,transform()让流式自动穿过链的每一层,而stream_events(version="v3")则以带反压的协议事件流支撑结构化输出与 Agent 执行可视化。理解这些机制之间的配合关系,是写出低延迟、高响应性、内存可控的 LangChain 应用的关键。相关实现与测试可进一步查阅 chat_models.py、ai.py、runnables/base.py 与 callbacks/base.py。

  • 人工智能
  • 大模型
  • AI Agent
  • Agent 框架
  • RAG

【免费下载链接】langchain

The agent engineering platform.

项目地址:https://gitcode.com/GitHub_Trending/la/langchain
点击查看免费下载

相关推荐

上一篇:Halo 页面布局契约落地实录:从 `layout :: html(...)` 约定到主题兼容状态机与 Console 提示
下一篇:用 InfoSpider 把 B 站观看历史和账号信息导出成本地 JSON 文件的完整流程

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/30 6:38:09

手把手在Switch大气层装好wiliwili:4步跑起来的B站客户端

手把手在Switch大气层装好wiliwili&#xff1a;4步跑起来的B站客户端 【免费下载链接】wiliwili 第三方B站客户端&#xff0c;目前可以运行在PC全平台、PSVita、PS4 、Xbox 和 Nintendo Switch上 项目地址: https://gitcode.com/GitHub_Trending/wi/wiliwili Switch大气…

作者头像 李华
网站建设 2026/9/30 6:37:57

Windows下VT控制权争夺:Hyper-V与安卓模拟器冲突真相

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/30 6:37:17

你的QQ空间历史说说找不到了?用GetQzonehistory一次全部找回

你的QQ空间历史说说找不到了&#xff1f;用GetQzonehistory一次全部找回 【免费下载链接】GetQzonehistory 获取QQ空间发布的历史说说 项目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory 想象一下&#xff0c;十年后你想重温大学时光&#xff0c;打开Q…

作者头像 李华