)
摘要本文介绍 LangChain 中 Agent 的流式输出机制重点讲解 updates、messages、custom 三种流式模式的区别与适用场景并演示如何通过 stream 函数、stream_mode 参数和 v2 版本结构获取实时输出最后给出一个完整的 custom 自定义数据流代码示例。内容参考于图灵AI大模型全栈流式输出它是一种让大模型边想边说的技术就是一个Agent它会运行很长的时间比如一个完成复杂的任务需要5分钟如果没有流式输出我们就需要在页面等待5分钟才能得到结果这样给人的体验就非常不好所以就出现了流式输出在LangChain中它提供了三种流式输出的模式如下图分别是 updates、messages、custom 三种模式updates它叫做流式状态更新它的核心是跟State状态打交道它就是把最新的状态发给我们就是说只要有人给状态的内容进行了更改它就会返回更改之后的状态基本上就是没执行完一个节点就会更新一次状态也就是返回一次最新的更改之后的状态我们最终看到的效果就是正在执行xx节点和执行完成和执行结果这些信息这样可以实时的掌握进度和变化它只会监控状态它的效果是整个任务的进度条它适合显示当前智能体在干嘛也就是说下图红框的messages被更新了它就可以得到messages它叫做大模型消息流现在的大模型它都是生成一个token就输出一个token这个只有在调用大模型的时候才会有看到的效果就是大模型一边想一边说它只会监控模型它的效果就是的我们在文档中打字的效果custom自定义数据流它是我们开发者自己实现的下方有示例它比较适合不支持流式输出的模型和函数内容自定义执行过程上方三种模式可以同时使用给 stream_mode 设置一个这样的[messages, updates, custom] 这样的值就可以全部获取到了下方示例中为了清晰是分开写的messages的效果就是一个如下图红框使用流式输出它就要通过 stream 函数向大模型进行提问注意下图红框它这里有一个版本想使用最新LangChain就要写v2如下图红框它默认v1版本这个v1和v2它们返回的效果和结构都不一样现在都推荐使用v2版本v2版本做了聚合如下图红框可以通过type取出不同的流式输出然后还需要传递一个模式如下图红框 stream_mode 的值用来设置模式然后通过type得到流式输出对象然后通过去data字段的值就可以得到最新的内容了如下图红框的代码custom自定义数据流示例# 从 langchain.agents 导入 create_agent构建 Agent 的高层工厂函数底层基于 LangGraph from langchain.agents import create_agent 从 langchain_openai 导入 ChatOpenAIOpenAI 兼容协议的客户端 这里通过 base_url 指向 DashScope通义千问的兼容端点 from langchain_openai import ChatOpenAI 从 langchain_core.utils.uuid 导入 uuid7生成 UUID v7时间有序的 UUID用于做线程 ID from langchain_core.utils.uuid import uuid7 从 langgraph.checkpoint.memory 导入 InMemorySaver LangGraph 的内存版 checkpointer按 thread_id 持久化对话状态让同一线程多次调用共享历史 from langgraph.checkpoint.memory import InMemorySaver 从 langgraph.config 导入 get_stream_writer 在工具/节点内部获取自定义流的写入器用于把自定义内容推送到 stream_modecustom 的流中 from langgraph.config import get_stream_writer 从 dotenv 导入 load_dotenv从 .env 文件加载环境变量到 os.environ from dotenv import load_dotenv 导入 os用于读取环境变量 import os 加载 .env 中的环境变量API Key、Base URL 等 load_dotenv() 初始化大模型 llm ChatOpenAI( modelqwen3.7-flash, # 模型名称 api_keyos.getenv(DASHSCOPE_API_KEY), # 从环境变量读取 DashScope 的 API Key base_urlos.getenv(DASHSCOPE_BASE_URL), # 从环境变量读取 DashScope 的 OpenAI 兼容端点 ) 定义一个工具LangChain 会读取函数名、类型注解、docstring 自动生成工具 schema LLM 通过这份 schema 决定何时、以什么参数调用它 def get_weather(city: str) - str: 获取指定城市的天气信息 # 在工具内部获取自定义流写入器只有 stream_modecustom 时下游才会消费 stream_writer get_stream_writer() # 向自定义流推送一条处理中的提示 stream_writer(f正在获取{city}的天气信息...) # 向自定义流推送一条完成的提示 stream_writer(f{city}的天气信息获取完成{city}的天气是晴朗的!) # 返回给 LLM 的工具调用结果会被包成 ToolMessage 塞回消息历史 return f{city}的天气是晴朗的! 创建 Agent 实例 agent create_agent( modelllm, # 使用的大模型 tools[get_weather], # 注册给 Agent 的工具列表 system_prompt你是一个乐于助人的助手能够回答用户关于天气的问题。请使用提供的工具来获取天气信息。, # 系统提示词约束角色和行为 checkpointerInMemorySaver(), # 内存版检查点按 thread_id 保存/恢复状态 ) 线程配置用 UUID v7 生成一个唯一 thread_id 同一个 config 在下面三次 stream 调用里复用所以三次共享同一份历史状态messages 会累加 config {configurable: {thread_id: str(uuid7())}} 流式模式 1updates 每次某个节点执行完毕、返回状态更新时把这一步的增量更新推出来粒度为节点级 print( 流式输出updates ) 调用 agent.stream 开启流式输出用 for 逐块消费 for chunk in agent.stream( # 参数 1input # 本次要合并进 State 的输入是增量不是完整 State # messages 字段会走父类 AgentState 绑定的 add_messages reducer → 追加而非覆盖 # {role: user, content: ...} 是 OpenAI 风格简写LangChain 会自动转成 HumanMessage {messages: [{role: user, content: 长沙天气怎么样?}]}, # 参数 2config # 运行时配置主要给 checkpointer 提供 thread_id用于定位/写入对应状态 configconfig, # 参数 3stream_mode # 指定流式输出的粒度/类型 # values → 每次推整个 State 快照 # updates → 每次推某个节点返回的增量更新本段用 # messages → 每次推一个 token/内容块下一段用 # custom → 推 get_stream_writer() 写入的内容被注释那段用 # debug → 推详细调试事件 stream_modeupdates, # 参数 4version # 输出数据格式版本v2 下每个 chunk 是统一的 {type: ..., data: ...} 结构 # 便于多种 stream_mode 混用时按 chunk[type] 分流 versionv2, ): # 判断 chunk 类型是不是 updates if chunk[type] updates: # chunk[data] 形如 {节点名: 该节点返回的状态更新} for step, data in chunk[data].items(): # 打印节点名 print(fstep: {step}) # 打印该节点返回的最后一条消息的内容块结构化内容 print(fcontent: {data[messages][-1].content_blocks}) 流式模式 2messages LLM 每产生一个 token/内容块就推出来一次粒度为 token 级适合做打字机效果 print( 流式输出messages ) 再次开启流式输出这次用 messages 模式 for chunk in agent.stream( # 参数 1input同上追加一条用户消息因为 config 相同会追加到历史 {messages: [{role: user, content: 长沙天气怎么样?}]}, # 参数 2config复用同一个 thread_id与上一段共享历史 configconfig, # 参数 3stream_modemessagestoken 级流式 # 每个 chunk 的 data 是 (消息块, metadata) 元组 stream_modemessages, # 参数 4versionv2统一的 {type: ..., data: ...} 结构 versionv2, ): # 判断 chunk 类型是不是 messages if chunk[type] messages: # 拆包token 是消息块对象metadata 是来源元信息 token, metadata chunk[data] # 打印该 token 来自哪个节点例如 model、tools print(fnode: {metadata[langgraph_node]}) # 打印内容块会输出很多空内容代表模型在思考阶段直到最后输出完整文本 print(fcontent: {token.content_blocks}) # 空行分隔便于阅读 print(\n) 流式模式 3custom 消费用户在工具/节点内通过 get_stream_writer() 写入的内容 本例里就是 get_weather 里那两句 正在获取.../...获取完成 注意不开启 custom 模式时stream_writer 写的内容不会被任何地方消费 print( 流式输出custom ) for chunk in agent.stream( # 参数 1input {messages: [{role: user, content: 长沙天气怎么样?}]}, # 参数 2config configconfig, # 参数 3stream_modecustom消费 get_stream_writer() 写入的内容 stream_modecustom, # 参数 4versionv2 versionv2, ): # 判断 chunk 类型是不是 custom if chunk[type] custom: # 直接打印 stream_writer 写入的原始内容 print(chunk[data])三种模式整合使用也就是 stream_mode[custom,messages,updates]这样的# 从 langchain.agents 导入 create_agent构建 Agent 的高层工厂函数底层基于 LangGraph from langchain.agents import create_agent # 从 langchain_openai 导入 ChatOpenAIOpenAI 兼容协议的客户端 # 这里通过 base_url 指向 DashScope通义千问的兼容端点 from langchain_openai import ChatOpenAI # 从 langchain_core.utils.uuid 导入 uuid7生成 UUID v7时间有序的 UUID用于做线程 ID from langchain_core.utils.uuid import uuid7 # 从 langgraph.checkpoint.memory 导入 InMemorySaver # LangGraph 的内存版 checkpointer按 thread_id 持久化对话状态让同一线程多次调用共享历史 from langgraph.checkpoint.memory import InMemorySaver # 从 langgraph.config 导入 get_stream_writer # 在工具/节点内部获取自定义流的写入器用于把自定义内容推送到 stream_modecustom 的流中 from langgraph.config import get_stream_writer # 从 dotenv 导入 load_dotenv从 .env 文件加载环境变量到 os.environ from dotenv import load_dotenv # 导入 os用于读取环境变量 import os # 加载 .env 中的环境变量API Key、Base URL 等 load_dotenv() # 初始化大模型 llm ChatOpenAI( modelqwen3.7-flash, # 模型名称 api_keyos.getenv(DASHSCOPE_API_KEY), # 从环境变量读取 DashScope 的 API Key base_urlos.getenv(DASHSCOPE_BASE_URL), # 从环境变量读取 DashScope 的 OpenAI 兼容端点 ) # 定义一个工具LangChain 会读取函数名、类型注解、docstring 自动生成工具 schema # LLM 通过这份 schema 决定何时、以什么参数调用它 def get_weather(city: str) - str: 获取指定城市的天气信息 # 在工具内部获取自定义流写入器只有 stream_modecustom 时下游才会消费 stream_writer get_stream_writer() # 向自定义流推送一条处理中的提示 stream_writer(f正在获取{city}的天气信息...) # 向自定义流推送一条完成的提示 stream_writer(f{city}的天气信息获取完成{city}的天气是晴朗的!) # 返回给 LLM 的工具调用结果会被包成 ToolMessage 塞回消息历史 return f{city}的天气是晴朗的! # 创建 Agent 实例 agent create_agent( modelllm, # 使用的大模型 tools[get_weather], # 注册给 Agent 的工具列表 system_prompt你是一个乐于助人的助手能够回答用户关于天气的问题。请使用提供的工具来获取天气信息。, # 系统提示词约束角色和行为 checkpointerInMemorySaver(), # 内存版检查点按 thread_id 保存/恢复状态 ) # 线程配置用 UUID v7 生成一个唯一 thread_id # 同一个 config 在下面三次 stream 调用里复用所以三次共享同一份历史状态messages 会累加 config {configurable: {thread_id: str(uuid7())}} print( 流式输出custom ) # noinspection PyTypeChecker for chunk in agent.stream( # 参数 1input {messages: [{role: user, content: 长沙天气怎么样?}]}, # 参数 2config configconfig, # 参数 3messages、custom、updates无顺序谁在前谁在后都可以 stream_mode[custom,messages,updates], # 参数 4versionv2 versionv2, ): # 判断 chunk 类型是不是 custom if chunk[type] custom: # 直接打印 stream_writer 写入的原始内容 print(chunk[data]) if chunk[type] updates: # 这里写 updates 模式的处理逻辑 print(chunk[data]) if chunk[type] messages: # 这里写 messages 模式的处理逻辑 print(chunk[data])