Agents 核心能力 [ 3 ]

人机协作(Human-in-the-loop)
接下来我们学习 Agent 当中如何实现人机协作能力。
首先说明核心实现方式:Agent 内置人机协作功能,直接使用 LangChain 官方预构建好的专用中间件即可完成开发。同时要区分开LangGraph 人机协作和Agent 人机协作的核心差异,二者人工介入的流程节点完全不同,我们分开对比讲解。
LangGraph 人机协作逻辑
LangGraph 的工作流由多个独立节点、连接线组成,执行链路灵活。 框架提供统一的中断方法,任意节点内都可以触发中断: 不管节点职责是调用大模型、处理数据、执行工具,只要在节点代码里写入中断逻辑,工作流就会立刻暂停,交由人工审核干预。 简单来说:LangGraph 可以在大模型执行前、任意工具执行前、数据处理前全流程任意位置打断流程,人工干预范围不受限制。
Agent 人机协作逻辑
Agent 的执行流程是固定固化的,完整链路:用户输入 → 大模型推理决策 → 调用工具 → 工具执行完毕返回结果 → 再次交给大模型整合输出。
结论:Agent 的人工审核中断仅能发生在工具调用阶段。 只有当大模型做出决策、准备调用某个工具时,流程才会暂停,交给人工判断是否放行本次工具调用; 我们无法在大模型推理执行前打断流程,不能干预大模型本身的执行,只能管控「工具能不能被调用」,这是二者最关键的区别。
Human-in-the-loop(HITL)为 AI Agent 添加人工审核与干预能力。HITL 中间件允许在 Agent 执行敏感工具调用前,暂停流程并等待人工审批。HITL 在自动化与风险控制之间建立平衡,适用于写文件、执行 SQL 等需要人工确认的操作。 其实现依赖 LangGraph 的持久化层(checkpointer),实现执行状态的保存与恢复。
三种中断决策类型
当 Agent 触发中断时,人工可做出以下三种响应(由策略配置决定哪些可用):
| 决策类型 | 说明 | 示例 |
|---|---|---|
✅ approve |
完全批准,工具按原参数执行 | 发送邮件草稿 |
✏️ edit |
允许修改工具参数后再执行 | 修改邮件收件人后发送 |
❌ reject |
拒绝执行,将反馈添加到对话中 | 拒绝 SQL 删除操作并说明原因 |
当中断触发、等待人工处理工具调用请求时,一共有三种可执行的干预方案,用来控制本次工具调用:
approve 完全批准:人工审核通过,工具使用原本的名称、原始参数直接执行,流程正常往下走。
reject 拒绝执行:直接驳回本次工具调用,同时可以填写拒绝原因,这条反馈消息会自动追加到对话历史里,交给大模型重新判断。
edit 修改参数后执行:工具本身允许调用,但传入的参数存在错误 / 不合理,人工可以修改工具入参,修改完成后再放行执行工具。
以上三种操作,就是 Agent 人机交互全部的干预手段,仅针对工具调用生效。
⚠️ 注意:当多个工具调用同时中断时,需按顺序逐一决策;编辑参数时请保守修改,避免影响模型后续判断。
配置与响应中断
实现该功能的核心中间件是 HumanInTheLoopMiddleware,创建 Agent 时直接注册到 middleware 列表即可。实例化中间件需要两个核心入参:
-
interrupt_on(必选):字典类型,用来定义每个工具是否触发中断、允许用户使用哪几种决策; -
description_prefix(可选):中断弹窗 / 日志的全局提示前缀,触发暂停时会先展示这段文字,例如配置"工具执行尚待批准"。
额外强制依赖:人机中断、流程恢复功能必须搭配 LangGraph 持久化检查点checkpointer使用,没有持久化存储无法保存中断状态、恢复会话。测试阶段可以使用内存存储InMemorySaver,生产环境需要使用数据库持久化方案(如 PostgresSaver)。
下面我们详细的对上述进行解释说明!
添加中间件与检查点
使用 HumanInTheLoopMiddleware 中间件完成人机协作,其配置选项主要包含以下两部分:
✅ 参数一:interrupt_on(必选)
一个字典,用于定义哪些工具需要触发人工中断以及允许哪些决策类型。 键为工具名称(字符串),值可以是以下三种形式:
| 值类型 | 含义 | 示例 |
|---|---|---|
True |
中断该工具,允许全部三种决策(approve、edit、reject) |
"工具1": True |
False |
不中断,自动批准(工具调用直接执行) | "工具2": False |
InterruptOnConfig 对象 |
精细控制:可指定允许的决策列表和自定义描述 | 见下方 |
InterruptOnConfig 对象的属性:
-
allowed_decisions:列表,可选值"approve"、"edit"、"reject"。决定人工可用的操作类型。 -
description:字符串或可调用函数,用于覆盖该工具的中断提示消息。若未指定,则使用全局description_prefix拼接而成。
我们以三类业务工具举例,讲解interrupt_on字典的三种配置写法,分别对应不同业务风险等级:
示例工具列表
write_file:写文件,会修改磁盘存储,属于高风险操作;execute_sql:执行 SQL 语句,可能删改数据库数据,中等风险;read_data:仅读取数据,无任何修改操作,低风险只读操作。
interrupt_on={
"write_file": True, # 全决策可用
"execute_sql": {"allowed_decisions": ["approve", "reject"]}, # 禁止编辑
"read_data": False, # 自动通过
}
1. 配置值 = False:自动放行,不触发中断
对应工具:read_data: False 含义:工具执行完全跳过人工审核,不会触发中断、不经过人机中间件,直接自动执行。 补充区分:和approve全部批准的表面效果一致,但底层流程不同:
-
False:完全绕过中间件,无中断逻辑; -
True:会经过中间件并触发中断,人工手动选择 approve 后放行。
2. 配置值 = True:全部三种决策开放
对应工具:write_file: True 含义:触发中断,人工三种操作全部可用:批准、拒绝、修改参数。适合高风险写文件、修改本地文件等操作。
3. 配置值 = InterruptOnConfig 字典:自定义开放部分决策
对应工具:execute_sql: {"allowed_decisions": ["approve", "edit"]} 通过allowed_decisions数组限定人工仅能使用指定操作,示例中仅允许批准、修改参数,禁用拒绝操作。 中断恢复时,用户无法执行不在列表内的决策,如果强行传入会直接报错。
配置总结三类场景
只读无风险工具:False,自动执行,跳过中断;
高风险操作,全权限管控:True,开放 approve/edit/reject 三种操作;
中等风险,限制人工操作:自定义allowed_decisions,仅开放部分审核权限。
✅ 参数二:description_prefix(可选)
字符串。作为中断消息的全局前缀,默认会在每个中断请求的描述前加上此文本。 最终消息格式:description_prefix + "\n\nTool: <tool_name>\nArgs: <arguments>" 工具级覆盖:若在 InterruptOnConfig 中指定了 description,则忽略全局前缀。
示例:
description_prefix="工具执行尚待批准"
中断时用户看到的消息开头为:
工具执行尚待批准
Tool: execute_sql
Args: {...}
通过合理配置 interrupt_on 和 description_prefix,开发者可以灵活控制哪些操作需要人工介入,以及如何向审核者展示信息。
完整代码执行流程拆解
-
调用
create_agent创建实例,传入模型、自定义工具列表、人机协作中间件、checkpointer 持久化存储; -
执行
agent.invoke()发起对话,必须传入包含thread_id的 config 配置,用来区分不同会话、保存中断状态; -
触发工具调用中断后,读取返回结果中的
interrupts字段,查看待人工审核的工具名称、参数、允许的操作类型; -
人工完成审核决策,使用
Command对象封装恢复指令,再次调用agent.invoke()传入同一thread_id的 config,继续推进工作流。
完整配置示例:
from langchain.agents import create_agent
from langchain.agents.middleware import HumanInTheLoopMiddleware
from langgraph.checkpoint.memory import InMemorySaver
agent = create_agent(
model="gpt-4.1",
tools=[write_file, execute_sql, read_data],
middleware=[
HumanInTheLoopMiddleware(
interrupt_on={
"write_file": True, # 允许全部三种决策
"execute_sql": {"allowed_decisions": ["approve", "reject"]}, # 禁止编辑
"read_data": False, # 自动批准,不中断
},
description_prefix="工具执行尚待批准", # 中断提示前缀
),
],
checkpointer=InMemorySaver(), # 生产环境请用持久化检查点(如PostgresSaver)
)
关键条件:
-
必须配置
checkpointer以支持中断恢复。 -
调用时需传入
thread_id以标识会话线程。 -
策略设计:对只读工具(如
read_data)可设False自动放行;对写操作建议至少配置approve和reject。
代码中 HumanInTheLoopMiddleware 配置表明:
-
write_file会中断,并允许批准、编辑、拒绝; -
execute_sql会中断,但仅允许批准或拒绝(不可编辑); -
read_data不中断,自动执行; -
中断提示将以 “工具执行尚待批准” 开头。
响应中断
前置版本说明
读取中断信息、使用 Command 恢复流程,需要指定调用参数version="v2",该版本规范从 LangChain 1.1 版本开始支持,优化了返回结果的数据结构,中断信息统一存放在result.interrupts,完整输出内容在result.v2.value。
触发中断后获取待审核动作
config = {"configurable": {"thread_id": "123"}}
result = agent.invoke(
{
"messages": [{"role": "user", "content": "删除数据库中的data表id=1的记录"}]
},
config=config, # 必须使用v2版本获取中断信息
)
print(result.interrupts)
输出包含 action_requests(工具名和参数)与 review_configs(允许的决策类型),结构示例:
Interrupt(
value={
"action_requests": [
{
"name": "execute_sql_tool",
"args": {"sql": "DELETE FROM data WHERE id = 1;"},
"description": "工具执行尚待批准\nTool: execute_sql_tool\nArgs: {'sql': 'DELETE FROM data WHERE id = 1;'}"
}
],
"review_configs": [
{
"action_name": "execute_sql_tool",
"allowed_decisions": ["approve", "reject"]
}
]
},
id="e678d1947c5028c88b98e961062d7b4d!"
)
提交决策继续执行
① approve 批准执行
拿到中断信息后,构造Command,decisions 数组内 type 设为"approve",传入同一个 thread_id 配置重新调用 agent,工具会使用原始参数正常执行。 执行完成后消息列表会新增ToolMessage,记录工具执行日志。
from langgraph.types import Command
# 批准执行
result = agent.invoke(
Command(resume={
"decisions": [{"type": "approve"}]
}),
config=config, # 同一thread_id
version="v2",
)
print(result.v2.value)
② reject 拒绝执行
仅允许配置了 reject 的工具才能使用该操作,decisions 内 type 设为"reject",新增message字段填写拒绝原因; 拒绝消息会封装为标准 ToolMessage 注入对话历史,大模型读取拒绝理由后重新规划后续动作,工具不会执行。 实操案例:对 write_file 写文件操作执行 reject,对话内会提示 “尝试写文件操作被驳回”,不会生成文件。
from langgraph.types import Command
# 拒绝执行
result = agent.invoke(
Command(resume={
"decisions": [
{
"type": "reject",
"message": "不批准,这是错误的,因为……。应该这样做……"
}
]
}),
config=config, # 同一thread_id
version="v2",
)
print(result.v2.value)
③ edit 修改参数后执行
decisions 内 type 设为"edit",同级新增edited_action字段,里面填写目标工具名称、修改后的新参数; 框架会替换本次工具调用的入参,使用修改后的参数执行工具。 实操案例:原始 SQL 为DELETE FROM data WHERE id = 1;,人工修改为id = 2,恢复后工具执行修改后的 SQL 语句。
from langgraph.types import Command
# 修改工具参数后再执行
result = agent.invoke(
Command(resume={
"decisions": [
{
"type": "edit",
# 带有工具名称和参数的编辑操作
"edited_action": {
# 要调整的工具名称,需与中断请求操作顺序保持一致
"name": "execute_sql_tool",
"args": {"key1": "new_value", "key2": "original_value"}
}
}
]
}),
config=config, # 同一thread_id
version="v2",
)
print(result.v2.value)
完整代码:
# ============================================================
# 1. 导入依赖
# ============================================================
from langchain.agents import create_agent
from langchain.agents.middleware import HumanInTheLoopMiddleware
from langchain.tools import tool
from langchain_deepseek import ChatDeepSeek
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import Command
# ============================================================
# 2. 定义工具(高风险 / 中风险 / 低风险)
# ============================================================
@tool
def write_file(filename: str, content: str) -> str:
"""
写入文件到磁盘(高风险操作)
"""
print(f"[文件写入] : {filename} -> [文件内容] : {content}")
return f"✅ 文件已写入: {filename}"
@tool
def execute_sql(sql: str) -> str:
"""
执行 SQL 语句(中风险操作)
"""
# 模拟执行 SQL
print(f"[SQL] 执行: {sql}")
return f"✅ SQL 执行成功: {sql}"
@tool
def read_data(query: str) -> str:
"""
读取数据(低风险只读操作)
"""
# 模拟读取数据
return f"📖 读取结果: 数据表包含 {len(query)} 条记录"
# ============================================================
# 3. 创建 Agent + 配置 HITL 中间件 + Checkpointer
# ============================================================
# 初始化模型
model = ChatDeepSeek(
model="deepseek-chat",
temperature=0.0,
)
# 创建内存检查点(生产环境请用 PostgresSaver)
checkpointer = InMemorySaver()
# 创建 Agent
agent = create_agent(
model=model,
tools=[write_file, execute_sql, read_data],
middleware=[
HumanInTheLoopMiddleware(
interrupt_on={
# 🔴 高风险:开放三种决策
"write_file": True,
# 🟡 中风险:只允许批准和拒绝,禁止编辑
"execute_sql": {"allowed_decisions": ["approve", "reject"]},
# 🟢 低风险:自动放行,不中断
"read_data": False,
},
description_prefix="⚠️ 工具执行需要人工审批",
)
],
checkpointer=checkpointer,
)
# ============================================================
# 4. 定义辅助函数:统一执行并处理中断
# ============================================================
def run_with_hitl(user_input: str, thread_id: str = "test_session"):
"""
执行 Agent,自动处理人机协作中断
"""
config = {"configurable": {"thread_id": thread_id}}
# 第一次调用:发起请求
result = agent.invoke(
{"messages": [{"role": "user", "content": user_input}]},
config=config,
version="v2", # 使用 v2 版本获取中断信息
)
# 检查是否有中断
if hasattr(result, "interrupts") and result.interrupts:
print("\n" + "=" * 60)
print("🚦 检测到人机中断,等待审批...")
print("=" * 60)
# 打印中断详情
interrupt_data = result.interrupts[0].value
print(f"\n📋 待审批操作:")
for action in interrupt_data["action_requests"]:
print(f" - 工具: {action['name']}")
print(f" - 参数: {action['args']}")
print(f" - 描述: {action.get('description', '无')}")
print(f"\n📌 允许的决策: {interrupt_data['review_configs']}")
return result, interrupt_data
print("\n✅ 无中断,执行完成")
print(f"结果: {result.value["messages"][-1].content}")
return result, None
# ============================================================
# 5. 三种决策响应函数
# ============================================================
def approve_execution(thread_id: str = "test_session"):
"""
决策1:批准执行(工具按原参数执行)
"""
config = {"configurable": {"thread_id": thread_id}}
result = agent.invoke(
Command(resume={
"decisions": [{"type": "approve"}]
}),
config=config,
version="v2",
)
print("\n✅ 已批准执行")
print(f"执行结果: {result.value["messages"][-1].content}")
return result
def reject_execution(reason: str, thread_id: str = "test_session"):
"""
决策2:拒绝执行(填写拒绝原因,反馈给 LLM)
"""
config = {"configurable": {"thread_id": thread_id}}
result = agent.invoke(
Command(resume={
"decisions": [
{
"type": "reject",
"message": reason,
}
]
}),
config=config,
version="v2",
)
print(f"\n❌ 已拒绝执行,原因: {reason}")
print(f"LLM 重新规划后的响应: {result.value["messages"][-1].content}")
return result
def edit_and_execute(tool_name: str, new_args: dict, thread_id: str = "test_session"):
"""
决策3:修改参数后执行
"""
config = {"configurable": {"thread_id": thread_id}}
result = agent.invoke(
Command(resume={
"decisions": [
{
"type": "edit",
"edited_action": {
"name": tool_name,
"args": new_args,
}
}
]
}),
config=config,
version="v2",
)
print(f"\n✏️ 已修改参数并执行")
print(f"修改后参数: {new_args}")
print(f"执行结果: {result.value["messages"][-1].content}")
return result
# ============================================================
# 6. 交互式 Demo:模拟完整人机协作流程
# ============================================================
def interactive_demo():
"""
交互式演示:用户选择不同的审批决策
"""
thread_id = "demo_session_001"
# 用户输入:触发高风险操作
user_input = "把 'Hello World' 写入到 test.txt 文件"
print("\n" + "=" * 60)
print(f"👤 用户输入: {user_input}")
print("=" * 60)
# 第一步:执行 Agent,触发中断
result, interrupt_data = run_with_hitl(user_input, thread_id)
if interrupt_data is None:
return
# 第二步:展示决策菜单
print("\n" + "=" * 60)
print("🎯 请选择审批决策:")
print("=" * 60)
print(" 1️⃣ approve - 批准执行(按原参数)")
print(" 2️⃣ reject - 拒绝执行(填写原因)")
print(" 3️⃣ edit - 修改参数后执行")
choice = input("\n请输入选项 (1/2/3): ").strip()
# 第三步:根据选择执行对应决策
if choice == "1":
approve_execution(thread_id)
elif choice == "2":
reason = input("请输入拒绝原因: ").strip()
reject_execution(reason, thread_id)
elif choice == "3":
tool_name = interrupt_data["action_requests"][0]["name"]
original_args = interrupt_data["action_requests"][0]["args"]
print(f"\n📌 原始参数: {original_args}")
# 让用户输入新参数(简单示例,实际可做结构化输入)
new_filename = input("请输入新文件名: ").strip()
new_content = input("请输入新内容: ").strip()
edit_and_execute(
tool_name=tool_name,
new_args={"filename": new_filename, "content": new_content},
thread_id=thread_id,
)
else:
print("⚠️ 无效选项")
# ============================================================
# 7. 手动场景测试(不需要交互)
# ============================================================
def manual_test_scenarios():
"""
手动测试三个场景
"""
thread_id = "manual_test"
# 场景1:高风险 - write_file 触发中断,批准执行
print("\n" + "🔵 场景1: 批准 write_file")
result, _ = run_with_hitl("把内容 'Hello' 写入 test.txt", thread_id)
if result and hasattr(result, "interrupts") and result.interrupts:
approve_execution(thread_id)
# 场景2:中风险 - execute_sql 触发中断,拒绝执行
print("\n" + "🔵 场景2: 拒绝 execute_sql")
result, _ = run_with_hitl("删除 users 表中 id=1 的记录", "thread_2")
if result and hasattr(result, "interrupts") and result.interrupts:
reject_execution("生产环境禁止删除操作,请先备份数据", "thread_2")
# 场景3:低风险 - read_data 自动放行,不触发中断
print("\n" + "🔵 场景3: read_data 自动执行(无中断)")
run_with_hitl("查询所有用户数据", "thread_3")
# 预期:无中断,直接返回结果
# ============================================================
# 8. 执行入口
# ============================================================
if __name__ == "__main__":
# 方式1:交互式演示
interactive_demo()
# 方式2:手动测试(取消注释运行)
# manual_test_scenarios()
Agent 人机协作仅能拦截、管控工具调用环节,无法在大模型推理前中断,和 LangGraph 全节点中断有本质区别;
通过HumanInTheLoopMiddleware实现,搭配interrupt_on对不同工具分级设置中断规则;
人工仅支持三种干预逻辑:批准、修改参数、拒绝,可通过allowed_decisions限制可用操作;
必须配置checkpointer持久化存储与thread_id会话标识,中断恢复依赖会话状态保存;
中断读取、流程恢复需要启用version="v2"规范,适配新版返回数据结构。
支持的流模式
接下来我们一起梳理 LangChain 里 Agent 所支持的全部流模式。
相信大家对流这个概念已经不陌生了,流式返回有多种能力:既能流式返回大语言模型输出的消息块;在 LangGraph 体系中,它还能流式返回图执行过程中走过的所有节点,同步带回节点更新后的完整状态;除此之外我们还能自定义数据反馈,这整套能力统一被称作流式系统,而 LangChain 里的 Agent 同样完整支持这套流式能力。
Agent 本身就是基于 LangGraph 封装实现的,所以它第一种流模式 updates,和我们之前在 LangGraph 里学过的 updates 逻辑、概念完全一致。它可以流式输出 Agent 当前执行到的阶段,同步返回该阶段完整的状态更新。除 updates 之外,还支持第二种流模式:流式返回大模型输出消息,也就是一个 token 逐次输出;第三种是自定义数据反馈模式;在这三种单独模式之外,我们还能把多种模式组合起来同步返回。下面我们逐个结合代码实操,完整走一遍编码流程。
LangChain Agent 流式系统能做什么?
-
流式传输代理进度:在每个代理步骤后获取状态更新。
-
流式传输 LLM token:实时生成语言模型产生的 token。
-
流式传输推理 token:输出模型的内部思考过程。
-
流式传输自定义更新:在代码中定义并发出用户自定义的信号(如 “已获取 10/100 条记录”)。
-
支持多种流模式:可根据需要选择
updates(代理进度)、messages(LLM 消息块)、custom(自定义数据)。
实操案例我们依旧沿用之前的天气查询工具,精简代码结构,只保留天气工具与 Agent 核心逻辑,删除所有无关中间件、多余导入包。
from langchain.agents import create_agent
from langchain_deepseek import ChatDeepSeek
from langchain.tools import tool
# ==================== 1. 定义 llm ====================
model = ChatDeepSeek(
model="deepseek-chat",
temperature=0.0,
)
# ==================== 2. 定义天气工具 ====================
@tool
def get_weather(city: str) -> str:
"""
查询指定城市的天气信息
Args:
city: 城市名称,如"北京"、"上海"
Returns:
天气信息字符串
"""
# 模拟天气API调用(实际项目可替换为真实API)
weather_data = {
"北京": "晴天,温度25°C,湿度40%",
"上海": "多云,温度28°C,湿度65%",
"广州": "雷阵雨,温度30°C,湿度80%",
"深圳": "晴转多云,温度27°C,湿度70%"
}
# 模拟API延迟
import time
time.sleep(0.5)
return weather_data.get(city, f"未找到{city}的天气信息,请检查城市名称")
# ==================== 3. 构建Agent ====================
system_prompt = """你是一个专业的天气助手,负责为用户提供准确、友好的天气信息。
## 你的职责
1. 当用户询问天气时,调用 get_weather 工具获取实时数据
2. 基于工具返回的数据,组织成清晰、易读的回复
3. 如果用户询问非天气相关问题,礼貌地引导回天气话题
## 回复规范
- 使用自然、友好的语气,像朋友聊天一样
- 结构化展示天气信息:温度、湿度、天气状况
- 如果有极端天气(暴雨、高温等),主动提醒注意事项
- 回复简洁明了,不要过度冗长
## 示例回复
用户问:"北京天气怎么样?"
回复:"北京今天晴天,温度25°C,湿度40%,非常适合户外活动!建议穿短袖,注意防晒哦~ 🌤️"
## 注意事项
- 不要编造天气数据,必须通过工具获取
- 如果工具返回错误,诚实告知用户并建议稍后重试
- 对于未支持的城市,友好提示并建议查询其他城市
"""
agent = create_agent(
model=model,
tools=[get_weather],
system_prompt=system_prompt,
)
# ==================== 5. 流式调用主程序 ====================
def main():
# 用户输入
user_input = "北京的天气如何"
print("=" * 50)
print(f"用户提问: {user_input}")
print("=" * 50)
print("\n开始流式输出...\n")
# ========== 方式一:使用v2版本(推荐) ==========
print("【v2 流式输出】")
print("-" * 40)
for chunk in agent.stream(
{"messages": [{"role": "user", "content": user_input}]},
stream_mode=["updates"], # 使用updates模式获取完整状态更新
version="v2" # 必须指定v2版本
):
print(chunk)
# v2统一格式:chunk是字典,包含type和data
# chunk_type = chunk.get("type")
# chunk_data = chunk.get("data")
#
# print(f"分片类型: {chunk_type}")
#
# # 解析data内容
# if chunk_type == "updates":
# # updates模式下,data包含节点更新信息
# for node_name, node_output in chunk_data.items():
# print(f" 节点: {node_name}")
#
# # 提取消息内容
# if "messages" in node_output:
# messages = node_output["messages"]
# for msg in messages:
# # 判断消息类型
# if hasattr(msg, "content") and msg.content:
# print(f" 内容: {msg.content[:200]}...") # 截断显示
# elif hasattr(msg, "tool_calls") and msg.tool_calls:
# tool_name = msg.tool_calls[0].get("name")
# tool_args = msg.tool_calls[0].get("args")
# print(f" 工具调用: {tool_name}({tool_args})")
#
# print("-" * 40)
# ========== 方式二:使用v1版本(对比) ==========
print("\n【v1 旧版流式输出】")
print("-" * 40)
for item in agent.stream(
{"messages": [{"role": "user", "content": user_input}]},
stream_mode=["updates"] # v1默认版本,不指定version
):
print(item)
# # v1格式:元组 (mode, data)
# print(f"模式: {mode}")
#
# # 解析chunk数据
# for node_name, node_output in chunk.items():
# print(f" 节点: {node_name}")
# if "messages" in node_output:
# messages = node_output["messages"]
# for msg in messages:
# if hasattr(msg, "content") and msg.content:
# print(f" 内容: {msg.content[:200]}...")
# elif hasattr(msg, "tool_calls") and msg.tool_calls:
# tool_name = msg.tool_calls[0].get("name")
# tool_args = msg.tool_calls[0].get("args")
# print(f" 工具调用: {tool_name}({tool_args})")
#
# print("-" * 40)
print("\n流式输出完成!")
# ==================== 7. 运行入口 ====================
if __name__ == "__main__":
main()
基础 Agent 写完后,常规同步调用我们会使用 .invoke() 方法,但流式场景不能用 invoke,需要改用 .stream() 流式调用方法。 .stream() 的第一个入参和 invoke 一致,都是传入用户输入,这里我们依旧使用提问「北京的天气如何」作为输入内容。
想要开启流式,我们需要额外配置流模式参数:
第一个入参是用户输入;
第二个参数用来指定 stream_mode 流模式; 如果我们需要逐阶段获取完整状态更新,就把 stream_mode 设置为 "updates";
同时必须指定 version="v2",这个 v2 版本要求 LangGraph 版本 ≥1.1,是新版统一返回格式规范。
.stream() 返回的是流式分片结果,不能用单个变量一次性接收,必须通过 for 循环遍历每一段 chunk 分片再处理打印。这里我们重点区分 v1(旧默认版)和 v2(新版)两种返回格式的差异。
v1 旧版本
v1 是当前默认版本,当同时使用多种组合流模式时,循环遍历拿到的是 (mode, chunk) 二元元组:
# 必须使用(mode, data)元组
for mode, chunk in agent.stream(
{"messages": [{"role": "user", "content": "上海今天天气怎么样?"}]},
stream_mode=["updates", "custom"],
):
print(mode) # "updates" or "custom"
print(chunk) # payload
mode:标识当前分片属于哪一种流模式(updates/custom);chunk:对应模式下的真实业务数据。
v2 新版本(推荐使用)(需要LangGraph>=1.1.)
v2 做了统一封装,循环中只会拿到单个 chunk 对象,不再需要元组解包:
# 统一格式—不再需要元组解包
for chunk in agent.stream(
{"messages": [{"role": "user", "content": "上海今天天气怎么样?"}]},
stream_mode=["updates", "custom"],
version="v2",
):
print(chunk["type"]) # "updates" or "custom"
print(chunk["data"]) # payload
chunk["type"]:等价于 v1 的 mode,标记当前分片模式;chunk["data"]:存放业务真实数据;- 额外携带
ns内置字段。
我们运行打印后能直观看到分片数据:本次天气案例会返回三段分片,分别对应三类节点执行:
========================================================================
用户提问:北京的天气如何
========================================================================
开始流式输出...
【v2 流式输出】
------------------------------
{'type': 'updates', 'ns': (), 'data': {'model': {'messages': [AIMessage(content='好的,我来查一下北京的天气情况!', additional_kwargs={'refusal'
{'type': 'updates', 'ns': (), 'data': {'tools': {'messages': [ToolMessage(content='晴天,温度25℃,湿度40%', name='get_weather', id='61aa438a
{'type': 'updates', 'ns': (), 'data': {'model': {'messages': [AIMessage(content='北京今天天气不错哦!☀️\n\n- **天气状况**:晴天\n- **温度**:25℃,
【v1 旧版流式输出】
------------------------------
('updates', {'model': {'messages': [AIMessage(content='好的,我来查一下北京的天气情况!', additional_kwargs={'refusal': None}, response_metadata:
('updates', {'tools': {'messages': [ToolMessage(content='晴天,温度25℃,湿度40%', name='get_weather', id='7d15136d-c949-4c42-9659-c370db1193
('updates', {'model': {'messages': [AIMessage(content='北京今天天气不错哦!☀️\n\n- **天气状况**:晴天\n- **温度**:25℃,体感舒适\n- **湿度**:40%,空
流式输出完成!
-
第一段:
model大模型节点,产出工具调用的 AI 消息; -
第二段:
tools工具节点,产出工具返回结果消息; -
第三段:
model大模型节点,产出整合后的最终回答消息。
这三段完整还原了 Agent 的执行链路,也就是 updates 模式流式输出的核心效果。
模式一:updates - 流式传输代理进度
在每个代理步骤(如一次 LLM 调用、一次工具执行)完成后,流式传输完整的状态更新。 示例:天气代理
for chunk in agent.stream(
{"messages": [{"role": "user", "content": user_input}]},
stream_mode=["updates"], # 使用updates模式获取完整状态更新
version="v2" # 必须指定v2版本
):
# print(chunk)
# v2统一格式:chunk是字典,包含type和data
chunk_type = chunk.get("type")
chunk_data = chunk.get("data")
print(f"分片类型: {chunk_type}")
# 解析data内容
if chunk_type == "updates":
# updates模式下,data包含节点更新信息
for node_name, node_output in chunk_data.items():
print(f" 节点: {node_name}")
# 提取消息内容
if "messages" in node_output:
messages = node_output["messages"]
for msg in messages:
# 判断消息类型
if hasattr(msg, "content") and msg.content:
print(f" 内容: {msg.content[:200]}...") # 截断显示
elif hasattr(msg, "tool_calls") and msg.tool_calls:
tool_name = msg.tool_calls[0].get("name")
tool_args = msg.tool_calls[0].get("args")
print(f" 工具调用: {tool_name}({tool_args})")
print("-" * 40)
输出:
chunk["data"] 内部是字典结构,键是执行阶段名(model / tools),值是该阶段更新后的完整状态,我们可以通过 for node_name, node_output in chunk_data.items() 遍历:
node_name:当前执行阶段名称;node_output:该阶段更新后的完整消息状态,包含所有对话、工具返回内容。
执行打印后能清晰看到完整执行流程:
-
用户提问「北京的天气如何」,首先进入
model节点,大模型生成调用天气工具的指令,此时状态里仅存在工具调用 AI 消息; -
接着进入
tools工具节点,执行天气工具,返回「北京阳光明媚」的结果; -
最后再次进入
model节点,大模型整合工具结果,生成最终回复文本。 这就是updates模式的完整作用:流式同步每一步执行后的更新状态。
==================================================
用户提问: 北京的天气如何
==================================================
开始流式输出...
【v2 流式输出】
----------------------------------------
分片类型: updates
节点: model
内容: 好的,我来查一下北京的天气情况!...
----------------------------------------
分片类型: updates
节点: tools
内容: 晴天,温度25°C,湿度40%...
----------------------------------------
分片类型: updates
节点: model
内容: 北京今天天气不错哦!☀️
- **天气状况**:晴天
- **温度**:25°C,体感舒适
- **湿度**:40%,空气比较干爽
非常适合出门活动!建议穿短袖或薄长袖,记得做好防晒哦~🌤️ 如果想了解其他城市的天气,随时问我!...
----------------------------------------
注意:将 version="v2"(需要 LangGraph >= 1.1.)传递给 stream() 或 astream() 以获取统一的输出格式。每个数据块都是一个具有 type、ns 和 data 键的 StreamPart 字典。
模式二:messages - 流式传输 LLM Token
从任何调用了LLM的节点中,流式传输token级别的增量消息块,实现类似 “逐字输出” 效果。每块是一个元组 (token, metadata)。
修改 stream_mode="messages",版本依旧保持 version="v2"。 循环内增加判断:仅当 chunk["type"] == "messages" 时处理数据。 chunk["data"] 是二元元组 (token, metadata):
-
第一个元素
token:模型输出的分片内容; -
第二个元素
metadata:元数据,记录分片来源节点信息。
从 token.content 可以取出分片文本打印,运行后能看到逐片输出效果:
-
第一阶段
model节点:逐分片输出调用工具的 JSON 参数,city: "北京"等内容会分段流式打印; -
中间
tools工具节点:一次性返回完整工具结果,不会分片流式输出; -
最后
model节点:逐分片输出整合后的最终回答文本。
我们可以通过元数据区分分片来源节点:metadata["langgraph_node"] 会标记当前分片来自 model 还是 tools 节点。 这里重点注意:messages 流模式只针对 LLM 大模型节点做分片流式传输,工具执行的返回内容不会拆分成 token 逐段输出。
示例:同上天气代理,使用 stream_mode="messages"
# 在循环外初始化一个字典来累积工具调用参数
tool_call_buffer = {}
current_tool_name = None
for chunk in agent.stream(
{"messages": [{"role": "user", "content": user_input}]},
stream_mode=["messages"],
version="v2"
):
chunk_type = chunk.get("type")
if chunk_type == "messages":
token, metadata = chunk["data"]
node_name = metadata.get("langgraph_node", "unknown")
if node_name == "model":
# 1. 处理完整的工具调用(第一个分片)
if hasattr(token, "tool_calls") and token.tool_calls:
valid_tool_calls = [tc for tc in token.tool_calls if tc.get("name")]
if valid_tool_calls:
print("\n🔧 [模型节点 - 工具调用]")
for tool_call in valid_tool_calls:
tool_name = tool_call.get("name")
tool_id = tool_call.get("id")
print(f" 📞 调用工具: {tool_name}")
# 初始化缓冲区
tool_call_buffer[tool_id] = {"name": tool_name, "args": ""}
current_tool_name = tool_name
# 2. 处理工具调用参数分片(通过 tool_call_chunks)
if hasattr(token, "tool_call_chunks") and token.tool_call_chunks:
for chunk_data in token.tool_call_chunks:
# 检查是否有参数分片
if chunk_data.get("args") is not None:
args_chunk = chunk_data["args"]
# 累积参数
if "args_buffer" not in locals():
args_buffer = ""
args_buffer += args_chunk
# 实时打印参数流(逐字符)
print(args_chunk, end="", flush=True)
# 3. 处理无效工具调用(也是参数分片的一种)
if hasattr(token, "invalid_tool_calls") and token.invalid_tool_calls:
for invalid_tc in token.invalid_tool_calls:
if invalid_tc.get("args"):
# 有些版本参数在 invalid_tool_calls 中
pass
# 4. 输出普通文本(非工具调用)
elif hasattr(token, "content") and token.content:
# 检查是否真的包含文本内容(排除工具调用的空content)
if not (hasattr(token, "tool_calls") and token.tool_calls):
print("\n💬 [模型节点 - 输出文本]")
print(token.content, end="", flush=True)
elif node_name == "tools":
print("\n🌤️ [工具节点 - 执行完成]")
if hasattr(token, "content") and token.content:
print(f" 📦 工具返回: {token.content}")
...
💬 [模型节点 - 输出文本]
天气
💬 [模型节点 - 输出文本]
情况
💬 [模型节点 - 输出文本]
!
🔧 [模型节点 - 工具调用]
📞 调用工具: get_weather
{"city": "北京"}
🌤️ [工具节点 - 执行完成]
📦 工具返回: 晴天,温度25°C,湿度40%
💬 [模型节点 - 输出文本]
北京
...
输出解读(部分):
- 你会看到
model节点逐块输出工具调用的 JSON 参数(如{'"', 'city', '":"', '北京', ...})。 - 然后
tools节点输出工具执行结果。 - 最后
model节点逐块输出最终回复文本(如'北京','的','天气'...)。
模式三:custom - 流式传输自定义更新
第三种 custom 自定义流模式,可以弥补工具无法主动推送中间进度的短板。如果我们需要在工具执行过程中,向外推送自定义提示、进度信息,就使用该模式。
与 LangGraph 用法一致,通过 get_stream_writer() 函数获取写入器,可以主动向流中推送任意自定义数据。 适用场景:报告工具执行进度(如 “正在查询数据库...”)、中间计算结果、调试信息等。 示例:在工具中发送自定义进度更新
具体使用步骤:
在工具函数内部导入 get_stream_writer() 方法,获取流写入器 writer;
调用 writer("自定义文本"),主动向流中推送自定义消息,例如「正在查询北京的天气数据」「北京天气查询完成」;
修改 stream_mode="custom",版本保持 v2;
循环内判断 chunk["type"] == "custom",直接打印 chunk["data"],就能收到工具推送的自定义文本。
运行后会单独打印我们在工具内定义的两句进度提示,实现工具执行过程的自定义反馈流式输出。
@tool
def get_weather(city: str) -> str:
"""
查询指定城市的天气信息
Args:
city: 城市名称,如"北京"、"上海"
Returns:
天气信息字符串
"""
writer = get_stream_writer()
writer(f"正在查询{city}的天气数据...")
# 模拟天气API调用(实际项目可替换为真实API)
weather_data = {
"北京": "晴天,温度25°C,湿度40%",
"上海": "多云,温度28°C,湿度65%",
"广州": "雷阵雨,温度30°C,湿度80%",
"深圳": "晴转多云,温度27°C,湿度70%"
}
# 模拟API延迟
import time
time.sleep(0.5)
writer(f"成功获取{city}的天气数据。")
return weather_data.get(city, f"未找到{city}的天气信息,请检查城市名称")
# ...
# ==================== 5. 流式调用主程序 ====================
def main():
# 用户输入
user_input = "北京的天气如何"
print("=" * 50)
print(f"用户提问: {user_input}")
print("=" * 50)
print("\n开始流式输出...\n")
# ========== 方式一:使用v2版本(推荐) ==========
print("【v2 流式输出】")
print("-" * 40)
# 在循环外初始化一个字典来累积工具调用参数
tool_call_buffer = {}
current_tool_name = None
for chunk in agent.stream(
{"messages": [{"role": "user", "content": user_input}]},
stream_mode=["custom"],
version="v2"
):
chunk_type = chunk.get("type")
chunk_data = chunk.get("data")
if chunk_type == "custom":
print(chunk_data)
输出:
==================================================
用户提问: 北京的天气如何
==================================================
开始流式输出...
【v2 流式输出】
----------------------------------------
正在查询北京的天气数据...
成功获取北京的天气数据。
模式四:组合模式 - 同时使用多种流模式
流模式支持组合使用,只需要把 stream_mode 传为列表,例如 stream_mode=["updates", "custom"],就能同时获取阶段状态更新、自定义推送消息两类分片。
依旧使用 v2 版本规范,循环遍历单个 chunk 对象:
-
通过
chunk_type判断当前分片是updates还是custom; -
通过
chunk_data取出对应模式的业务数据。
for chunk in agent.stream(
{"messages": [{"role": "user", "content": user_input}]},
stream_mode=["updates", "custom"],
version="v2"
):
chunk_type = chunk.get("type")
chunk_data = chunk.get("data")
print(f"流模式: {chunk_type}")
print(f"内容: {chunk_data}\n")
==================================================
用户提问: 北京的天气如何
==================================================
开始流式输出...
【v2 流式输出】
----------------------------------------
流模式: updates
内容: {'model': {'messages': [AIMessage(content='好的,我来查一下北京的天气情况!', ...
流模式: custom
内容: 正在查询北京的天气数据...
流模式: custom
内容: 成功获取北京的天气数据。
流模式: updates
内容: {'tools': {'messages': [ToolMessage(content='晴天,温度25°C,湿度40%', name='get_weather',...
流模式: updates
内容: {'model': {'messages': [AIMessage(content='北京今天天气不错哦!☀️\n\n- **天气状况**:晴天\n- **温度**:25°C,体感舒适\n- **湿度**:40%,空气比较干爽\n\n非常适合出门活动!建议穿短袖或薄长袖,记得做好防晒哦~🌤️ 这么好的天气,出去走走心情也会变好呢!😊', ...
执行时分片会交替返回:
-
先返回
type="updates":大模型生成工具调用指令的阶段状态; -
进入工具执行时,交替返回
type="custom":工具内自定义的进度提示; -
工具执行完毕,再次返回
type="updates":工具执行完成后的状态; -
最终返回
type="updates":大模型整合回答后的完整状态。
补充区分 v1 /v2 在组合模式下的写法差异:
-
v1 组合模式:循环接收
(mode, chunk)元组,先判断 mode 再取数据; -
v2 组合模式:仅接收单个 chunk 对象,依靠
chunk["type"]区分模式,结构更统一简洁。
将模式作为列表传入,如 stream_mode=["updates", "custom"]。每次迭代返回一个包含 type(标识模式)和 data(对应负载)的字典。 示例:同时使用 updates 和 custom 模式
模式实践
接下来我们讲解流式能力几个高频落地实践场景,在掌握流模式基础用法后,我们梳理它在业务里最常用的落地场景。第一个核心场景就是流式传输模型推理、思考 Token。
流式传输推理 / 思考 Token
某些模型(如 Claude、OpenAI 系列)在生成最终答案前会进行内部推理。将这些推理内容实时输出,可以增强应用的透明度与可解释性。
我们可以通过配置开启推理输出,把模型幕后的思考步骤实时展示给用户。 举个直观例子:我们向模型提问 “你好”,交互界面里会先出现一段 “正在思考” 的文字,这段内容就是大模型原生输出的推理过程,我们可以通过代码配置把这段思考内容完整捕获并流式打印。
实现原理:LangChain 将不同提供商的推理内容(Anthropic 的thinking块、OpenAI 的reasoning摘要等)统一规范为 content_blocks 中 type: "reasoning" 的块。
代码示例:
以 OpenAI 系列模型为例,想要开启推理输出,初始化 ChatOpenAI 时需要新增 reasoning 配置项,内部包含两个核心参数:
- effort(推理程度):可选
low、medium、high三档,控制模型思考的深度; - summary(推理摘要):支持
detailed(详细完整推理)、auto(自动精简)、None(不生成摘要)三种模式。
完成模型推理开关配置后,我们结合 Agent 的流式调用,就能实现推理过程逐 Token 实时流式输出的效果。
以 DeepSeek 系列模型为例,想要开启推理输出,初始化 ChatOpenAI 时需要新增 extra_body 配置项,内部包含推理参数:
核心参数
reasoning_effort(推理程度):可选 low、medium、high 三档,控制模型思考的深度;
⚠️ 注意:DeepSeek 目前不支持
summary(推理摘要)参数,OpenAI 的summary功能在 DeepSeek 中暂不可用。
配置示例
# DeepSeek 推理配置
model = ChatDeepSeek(
model="deepseek-reasoner", # 或 "deepseek-chat"
temperature=0.7,
streaming=True, # 必须开启流式
# 推理参数通过 extra_body 传递
extra_body={
"thinking": {"type": "enabled"}, # 开启思考模式
"reasoning_effort": "high", # 推理程度:low/medium/high
},
max_tokens=4096,
)
基础代码骨架复用:代码整体结构和之前天气查询 Agent 基本一致,仅做两处删减调整:
-
保留
stream_mode="messages":想要捕获推理分片,必须使用messages流模式; -
移除多余自定义工具相关冗余代码,精简逻辑;
-
Agent 绑定时,传入我们刚刚配置好、开启推理能力的
model对象。
from langchain.agents import create_agent
from langchain_deepseek import ChatDeepSeek
from langchain.tools import tool
# ==================== 1. 配置模型,启用推理输出 ====================
# model = ChatOpenAI(
# model="gpt-5-mini",
# reasoning={ # 关键配置
# "effort": "medium", # 推理程度: 'low', 'medium', or 'high'
# "summary": "auto", # 推理摘要: 'detailed', 'auto', or None
# }
# )
model = ChatDeepSeek(
model="deepseek-chat",
temperature=0.0,
streaming=True,
extra_body={
"thinking": {"type": "enabled"},
"reasoning_effort": "high",
},
max_tokens=4096,
)
# ==================== 2. 定义天气工具 ====================
@tool
def get_weather(city: str) -> str:
"""
查询指定城市的天气信息
Args:
city: 城市名称,如"北京"、"上海"
Returns:
天气信息字符串
"""
# 模拟天气API调用(实际项目可替换为真实API)
weather_data = {
"北京": "晴天,温度25°C,湿度40%",
"上海": "多云,温度28°C,湿度65%",
"广州": "雷阵雨,温度30°C,湿度80%",
"深圳": "晴转多云,温度27°C,湿度70%"
}
# 模拟API延迟
import time
time.sleep(0.5)
return weather_data.get(city, f"未找到{city}的天气信息,请检查城市名称")
# ==================== 3. 构建Agent ====================
agent = create_agent(
model=model,
tools=[get_weather],
)
流式循环基础结构:流式调用依旧使用 .stream() 方法,指定 version="v2" 新版返回格式,传入用户提问「北京的天气如何?」; 循环遍历每一段 chunk,判断 chunk["type"] == "messages" 后,从 chunk["data"] 解包二元元组 (token, metadata)。
解析 Token 中的推理内容核心逻辑:token 对象里的 content_blocks 是存储所有输出分片的核心字段,它是一个列表,列表内每一个字典分片依靠 type 字段区分内容类型:
-
type="text":最终给用户展示的正文回答; -
type="reasoning":模型内部推理、思考过程。
我们通过列表推导式分别筛选两类内容:
-
新建
reasoning列表,过滤所有type="reasoning"的分片,存放思考过程; -
新建
text列表,过滤所有type="text"的分片,存放最终回答正文。
分情况打印输出
-
推理内容:先判断
reasoning列表非空,取出列表第 0 个分片里的reasoning字段内容,使用end=""实现无换行逐字流式打印; -
正文内容:判断
text列表非空,取出分片内text字段逐段输出。
# ==================== 4. 流式调用主程序 ====================
def main():
# 使用 stream_mode = "messages" 并过滤 reasoning 块
for token, metadata in agent.stream(
{"messages": [{"role": "user", "content": "上海的天气如何?"}]},
stream_mode="messages"
):
if not isinstance(token, AIMessageChunk):
continue
# 推理内容
reasoning = [b for b in token.content_blocks if b["type"] == "reasoning"]
# 响应内容
text = [b for b in token.content_blocks if b["type"] == "text"]
if reasoning:
print(f"[reasoning] {reasoning[0]['reasoning']}", end="")
if text:
print(f"[text] {text[0]['text']}", end="")
[reasoning] [reasoning] 用户[reasoning] 想知道[reasoning] 上海的[reasoning] 天气[reasoning] 。[reasoning] 让我[reasoning] 查询[reasoning] 一下[reasoning] 。[reasoning] [reasoning] 好的[reasoning] ,[reasoning] 我已经[reasoning] 获取[reasoning] 到[reasoning] 上海的[reasoning] 天气[reasoning] 信息[reasoning] 了[reasoning] 。[reasoning] 让我[reasoning] 用[reasoning] 中文[reasoning] 回复[reasoning] 用户[reasoning] 。[text] 上海的[text] 天气[text] 情况[text] 如下[text] :
[text] -[text] **[text] 天气[text] 状况[text] **[text] :[text] 多云[text] ⛅[text]
[text] -[text] **[text] 温度[text] **[text] :[text] 28[text] °[text] C[text]
[text] -[text] **[text] 湿度[text] **[text] :[text] 65[text] %
[text] 总体来说[text] ,[text] 上海[text] 今天[text] 天气[text] 还不错[text] ,[text] 多云[text] 间[text] 晴[text] ,[text] 温度[text] 适宜[text] ,[text] 比较[text] 舒适[text] 。[text] 如果[text] 外出[text] 的话[text] 记得[text] 带[text] 把[text] 伞[text] ,[text] 以防[text] 万一[text] 哦[text] ~
关键点
-
必须确保模型支持并开启了推理输出(如
reasoning参数)。注意:模型不同,参数不同(Anthropic 系列参数为thinking) -
无论使用哪个提供商,都可通过
content_blocks中的type: "reasoning"统一访问推理内容。 -
通常配合
stream_mode="messages"使用,逐块输出推理文本。
流式传输工具调用
接下来我们继续讲解流模式的第二个落地实践场景 —— 流式传输工具调用。
我们之前使用messages流模式时,只能流式输出模型最终的正文回答;但大模型在判断需要调用工具的阶段,同样会产生分片流式输出。 这个阶段的输出节点依旧是model大模型节点,只是返回内容不再是普通文本,而是调用工具的决策信息。对应数据里的tool_call_chunks字段,工具调用所需的参数会一段一段分片流式输出。
也就是说在代理调用工具时,我们可能希望:
-
实时显示工具调用的参数 JSON 片段(增量)。
-
最终获取完整解析后的工具调用消息(用于后续逻辑)。
所以,如果我们想要完整捕捉、展示整套工具调用流程,单一messages模式不足以覆盖全部信息,需要组合流模式使用。
-
搭配
stream_mode="messages":用来实时流式输出工具选择、工具参数的增量分片; -
搭配
stream_mode="updates":分阶段输出节点执行完成后的完整状态,model节点执行完毕后,会返回完整、解析好的tool_calls工具调用信息,tools工具节点执行完成后,会返回工具执行后的完整结果状态。
两种模式组合,就能完整实现工具调用全流程的流式传输。
具体的实现原理:
-
stream_mode="messages":提供AIMessageChunk,其中包含tool_call_chunks(增量参数)。 -
stream_mode="updates":在步骤(如model节点)完成后,提供完整的AIMessage,其中包含已解析的tool_calls。
代码示例(组合模式):
代码主体沿用上面的天气查询 Agent,其中两个辅助函数是用来格式化打印输出的,核心逻辑集中在流式循环部分:
【两个辅助函数】
# ==================== 3. 辅助函数 ====================
def render_chunk(token: AIMessageChunk):
if token.text:
print(token.text, end="|")
if token.tool_call_chunks:
print(token.tool_call_chunks) # 增量块
def render_completed(msg: AnyMessage):
if isinstance(msg, AIMessage) and msg.tool_calls:
print(f"完整工具调用: {msg.tool_calls}")
if isinstance(msg, ToolMessage):
print(f"工具响应: {msg.content_blocks}")
版本选择:组合模式必须使用version="v2"新版格式,循环仅接收单个chunk对象,不需要像 v1 版本那样解包元组;
用户输入:提问依旧使用「上海的天气如何?」;
双类型分片处理:组合模式下会收到两类chunk,代码通过if/elif分别处理:
-
chunk["type"] == "messages":大模型实时增量分片数据; -
chunk["type"] == "updates":每个节点执行完成后的完整阶段状态。
1. messages 增量分片处理逻辑
从chunk["data"]解包得到二元元组(token, metadata),代码里增加过滤逻辑:只处理来自model大模型节点的AIMessageChunk消息,过滤掉tools工具节点一次性返回的完整结果。
拿到token对象后,分别处理两类内容:
-
普通文本正文:读取
token.text逐段输出; -
工具调用增量分片:读取
token.tool_call_chunks字段,这个字段是列表格式,存放分片的 JSON 工具参数,能实现参数逐段流式打印。
输出顺序有固定逻辑:先流式输出工具调用参数分片,全部输出完成后,再流式输出模型最终正文文本。我们可以把这套增量参数流式输出交给前端渲染,前端就能实时展示模型正在拼接工具参数的全过程。
2. updates 完整阶段状态处理逻辑
遍历chunk["data"]内的所有节点更新,仅处理model、tools两类核心节点:
-
如果更新消息是
AIMessage且携带完整tool_calls,打印完整、解析完毕的工具调用信息; -
如果更新消息是
ToolMessage,打印工具执行完成后的完整响应内容。
这部分相当于全流程日志输出,能完整记录两条链路:
-
增量链路:通过
messages实时流式打印碎片化工具调用参数; -
阶段完整链路:通过
updates在每个节点结束后,打印完整工具调用、完整工具返回结果。
这套日志数据可以直接透传给前端,前端可自由设计展示样式,也能和模型推理思考过程合并展示。
# ==================== 4. 流式调用主程序 ====================
def main():
for chunk in agent.stream(
{"messages": [{"role": "user", "content": "上海的天气如何?"}]},
stream_mode=["messages", "updates"],
version="v2"
):
chunk_type = chunk["type"]
chunk_date = chunk["data"]
if chunk_type == "messages":
token, metadata = chunk_date
if isinstance(token, AIMessageChunk):
render_chunk(token)
elif chunk["type"] == "updates":
for source, update in chunk["data"].items():
if source in ("model", "tools"): # 关注模型和工具节点
render_completed(update["messages"][-1])
[{'name': 'get_weather', 'args': '', 'id': 'call_00_Ywnze68uUURfgDJGfoZr1140', 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '{', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': 'city', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': ': ', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '上海', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '}', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
完整工具调用: [{'name': 'get_weather', 'args': {'city': '上海'}, 'id': 'call_00_Ywnze68uUURfgDJGfoZr1140', 'type': 'tool_call'}]
工具响应: [{'type': 'text', 'text': '多云,温度28°C,湿度65%'}]
上海|当前|天气|情况|如下|:
|-| **|天气|状况|**|:|多云| 🌤|️|
|-| **|温度|**|:|28|°|C|
|-| **|湿度|**|:|65|%
|整体|来看|天气|还算|不错|,|多云|天气|不会|太|晒|,|温度|也比较|舒适|。|不过|湿度|稍|高|,|体|感|可能会|稍微|有些|闷|热|,|出门|建议|注意|补水|哦|!|
组合messages+updates流模式后,我们能完整流式还原工具调用全流程:
-
实时分片流式输出工具调用的 JSON 参数;
-
模型节点执行完毕后,输出完整解析后的工具调用记录;
-
工具执行完成后,输出工具返回的完整结果;
-
最终流式输出大模型整合后的回答文本。
关键点
-
tool_call_chunks是增量数据,可用于实时显示(如 “正在输入参数…”)。 -
完整的
tool_calls需从updates模式中获取,通常在model节点完成后出现。 -
如果消息未保存在状态中,可通过累积
AIMessageChunk(使用+操作符)来重构完整消息。
流式传输与人机协作 (Human-in-the-Loop)
第三个核心实践场景:流式传输人机协作交互。当 Agent 执行到关键操作(调用工具)时,流程会主动暂停,等待人工审批、修改工具参数后,再恢复执行;整个中断、人工决策、恢复执行的全流程,都可以通过stream流式输出全部信息。
实现原理:
-
使用
HumanInTheLoopMiddleware中间件,指定需要中断的工具。 -
流式传输时,捕获
updates模式中的__interrupt__节点。 -
构造
Command(resume=...)来恢复执行。
代码无需从零编写,核心配置要点如下:
Agent 初始化时挂载HumanInTheLoopMiddleware人机交互中间件,配置interrupt_on={"工具名": True},开启指定工具的中断拦截;
支持三种人工决策策略,恢复执行时可任选:
-
approve:审批通过,直接执行原工具参数;
-
edit:修改工具入参后再执行;
-
reject:拒绝本次工具调用;
必须配置thread_id会话标识,用来持久化会话状态,保存中断现场。
# ==================== 3. 构建Agent ====================
agent = create_agent(
model=model,
tools=[get_weather],
checkpointer=InMemorySaver(),
middleware=[HumanInTheLoopMiddleware(
interrupt_on={
"get_weather": True,
},
)],
)
流式调用时指定stream_mode="updates",分阶段捕获节点状态:
-
遍历
chunk["data"]中的节点更新,识别特殊节点__interrupt__中断节点; -
将所有中断信息存入列表统一管理;
-
解析中断数据
interrupt.value,里面包含被拦截工具的名称、原始入参、中断提示描述,直接打印展示给用户,告知当前哪一步工具调用需要人工处理。
def render_interrupt(interrupt: Interrupt) -> None:
for request in interrupt.value["action_requests"]:
print(f"需要人工审批: {request}")
# ==================== 4. 流式调用主程序 ====================
def main():
config = {"configurable": {"thread_id": "1"}}
interrupts = []
# 第一次流式: 遇到中断
for chunk in agent.stream(
{"messages": [{"role": "user", "content": "查询北京和上海的天气"}]},
config=config,
stream_mode="updates",
version="v2",
):
if chunk["type"] == "updates":
for source, update in chunk["data"].items():
if source == "__interrupt__":
# 解析中断,展示给用户,收集决策
interrupts.extend(update[0])
render_interrupt(update[0])
本次演示提问为「查询北京和上海的天气」,会连续两次触发天气工具中断,控制台会先后打印两条审批提示,分别对应北京、上海两次工具调用。
需要人工审批: {'name': 'get_weather', 'args': {'city': '北京'}, 'description': "Tool execution requires approval\n\nTool: get_weather\nArgs: {'city': '北京'}"}
需要人工审批: {'name': 'get_weather', 'args': {'city': '上海'}, 'description': "Tool execution requires approval\n\nTool: get_weather\nArgs: {'city': '上海'}"}
下面我们进行人工决策与恢复执行
1. 构造决策指令
收集到全部中断后,构造决策列表decisions,每条决策包含:
-
interrupt_id:对应中断的唯一标识,指定该决策作用于哪一次中断; -
action:人工操作行为,包含审批 / 编辑配置; 本次演示决策规则:
第一次北京天气调用:approve直接审批通过;
第二次上海天气调用:edit编辑参数,将城市从「上海」修改为「西安」。
2. 恢复执行流式处理
调用agent.stream()并传入Command(resume=decisions)恢复流程:
# ==================== 4. 流式调用主程序 ====================
def main():
config = {"configurable": {"thread_id": "1"}}
interrupts = []
# 第一次流式:完整遍历所有chunk,收集全部中断后再处理
print("===== 首次执行,捕获工具中断 =====")
stream_iter = agent.stream(
{"messages": [{"role": "user", "content": "查询北京和上海的天气"}]},
config=config,
stream_mode="updates",
version="v2",
)
# 完整遍历所有输出块,等待全部中断收集完毕
for chunk in stream_iter:
print(f"Chunk: {chunk}") # 调试输出
if chunk["type"] == "updates":
for source, update in chunk["data"].items():
if source == "__interrupt__":
# update 是一个元组,包含 Interrupt 对象
if isinstance(update, tuple):
for interrupt in update:
if isinstance(interrupt, Interrupt):
interrupts.append(interrupt)
render_interrupt(interrupt)
elif isinstance(update, list):
for interrupt in update:
if isinstance(interrupt, Interrupt):
interrupts.append(interrupt)
render_interrupt(interrupt)
else:
# 单个 Interrupt 对象
if isinstance(interrupt, Interrupt):
interrupts.append(interrupt)
render_interrupt(interrupt)
# 关键:全部chunk遍历完成后,再读取中断列表
print(f"\n共捕获 {len(interrupts)} 个人工审批中断")
for i, interrupt in enumerate(interrupts):
print(f"中断 {i}: {interrupt.value}")
# 🔥 关键修复:为每个中断创建决策
decisions = []
for interrupt in interrupts:
# 从 interrupt.value 中提取 action_requests
if isinstance(interrupt.value, dict) and "action_requests" in interrupt.value:
action_requests = interrupt.value["action_requests"]
# 为每个 action_request 创建决策
for action_req in action_requests:
# 如果是北京,直接批准
if action_req["args"]["city"] == "北京":
decisions.append({
"type": "approve",
})
# 如果是上海,修改为西安
elif action_req["args"]["city"] == "上海":
decisions.append({
"type": "edit",
"edited_action": {
"name": "get_weather",
"args": {"city": "西安"},
},
})
else:
# 其他情况默认批准
decisions.append({
"type": "approve",
})
else:
# 如果是单个工具调用(兼容旧格式)
decisions.append({
"type": "approve",
})
print(f"\n生成 {len(decisions)} 个审批决策")
# 🔥 关键修复:构造正确的 resume 格式
# 对于包含多个 action_requests 的中断,需要返回一个包含 'decisions' 键的字典
resume_data = {
"decisions": decisions
}
# 第二次流式:恢复执行
print("\n===== 提交审批决策,恢复Agent执行 =====")
for chunk in agent.stream(
Command(resume=resume_data),
config=config,
stream_mode=["messages", "updates"],
version="v2",
):
if chunk["type"] == "messages":
token, meta = chunk["data"]
if isinstance(token, AIMessageChunk):
render_chunk(token)
elif chunk["type"] == "updates":
for source, update in chunk["data"].items():
if source in ("model", "tools"):
render_completed(update["messages"][-1])
elif source == "__interrupt__":
# 处理恢复后可能的新中断
if isinstance(update, tuple):
for interrupt in update:
if isinstance(interrupt, Interrupt):
render_interrupt(interrupt)
elif isinstance(update, Interrupt):
render_interrupt(update)
===== 首次执行,捕获工具中断 =====
Chunk: {'type': 'updates', 'ns': (), 'data': {'model': {'messages': [AIMessage(content='好的,我来同时查询北京和上海的天气信息。', additional_kwargs={'reasoning_content': '用户想查询北京和上海的天气,我需要同时调用两次get_weather工具。'}, response_metadata={'finish_reason': 'tool_calls', 'model_name': 'deepseek-v4-flash', 'system_fingerprint': 'fp_8b330d02d0_prod0820_fp8_kvcache_20260402', 'model_provider': 'deepseek'}, id='lc_run--019f0f63-8db3-7aa2-93c7-a77fb32b8099', tool_calls=[{'name': 'get_weather', 'args': {'city': '北京'}, 'id': 'call_00_ANG6pMpLK3gZtfs5yVWd7089', 'type': 'tool_call'}, {'name': 'get_weather', 'args': {'city': '上海'}, 'id': 'call_01_eoBRiOH2gy3nCdVVziEu5707', 'type': 'tool_call'}], invalid_tool_calls=[], usage_metadata={'input_tokens': 300, 'output_tokens': 103, 'total_tokens': 403, 'input_token_details': {'cache_read': 256}, 'output_token_details': {'reasoning': 17}})]}}}
Chunk: {'type': 'updates', 'ns': (), 'data': {'__interrupt__': (Interrupt(value={'action_requests': [{'name': 'get_weather', 'args': {'city': '北京'}, 'description': "Tool execution requires approval\n\nTool: get_weather\nArgs: {'city': '北京'}"}, {'name': 'get_weather', 'args': {'city': '上海'}, 'description': "Tool execution requires approval\n\nTool: get_weather\nArgs: {'city': '上海'}"}], 'review_configs': [{'action_name': 'get_weather', 'allowed_decisions': ['approve', 'edit', 'reject', 'respond']}, {'action_name': 'get_weather', 'allowed_decisions': ['approve', 'edit', 'reject', 'respond']}]}, id='dcfb2ffcb38e7c328df6703b291d05ca'),)}}
需要人工审批: {'action_requests': [{'name': 'get_weather', 'args': {'city': '北京'}, 'description': "Tool execution requires approval\n\nTool: get_weather\nArgs: {'city': '北京'}"}, {'name': 'get_weather', 'args': {'city': '上海'}, 'description': "Tool execution requires approval\n\nTool: get_weather\nArgs: {'city': '上海'}"}], 'review_configs': [{'action_name': 'get_weather', 'allowed_decisions': ['approve', 'edit', 'reject', 'respond']}, {'action_name': 'get_weather', 'allowed_decisions': ['approve', 'edit', 'reject', 'respond']}]}
共捕获 1 个人工审批中断
中断 0: {'action_requests': [{'name': 'get_weather', 'args': {'city': '北京'}, 'description': "Tool execution requires approval\n\nTool: get_weather\nArgs: {'city': '北京'}"}, {'name': 'get_weather', 'args': {'city': '上海'}, 'description': "Tool execution requires approval\n\nTool: get_weather\nArgs: {'city': '上海'}"}], 'review_configs': [{'action_name': 'get_weather', 'allowed_decisions': ['approve', 'edit', 'reject', 'respond']}, {'action_name': 'get_weather', 'allowed_decisions': ['approve', 'edit', 'reject', 'respond']}]}
生成 2 个审批决策
===== 提交审批决策,恢复Agent执行 =====
工具响应: 晴天,温度25°C,湿度40%
工具响应: 未找到西安的天气信息,请检查城市名称
抱歉|,|我|查|错了|城市|,|让我|重新|查询|上海的|天气|。|[{'name': 'get_weather', 'args': '', 'id': 'call_00_d2N4YxiQ2POPWBIKy9W58485', 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '{', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': 'city', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': ': ', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '上海', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '"', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
[{'name': None, 'args': '}', 'id': None, 'index': 0, 'type': 'tool_call_chunk'}]
完整工具调用: [{'name': 'get_weather', 'args': {'city': '上海'}, 'id': 'call_00_d2N4YxiQ2POPWBIKy9W58485', 'type': 'tool_call'}]
需要人工审批: {'action_requests': [{'name': 'get_weather', 'args': {'city': '上海'}, 'description': "Tool execution requires approval\n\nTool: get_weather\nArgs: {'city': '上海'}"}], 'review_configs': [{'action_name': 'get_weather', 'allowed_decisions': ['approve', 'edit', 'reject', 'respond']}]}
中断阶段输出:控制台打印两次工具审批提示,告知用户北京、上海天气调用需要人工确认;
恢复执行初始输出:直接执行修改后的工具,打印工具响应:北京天气晴朗、西安天气晴朗;
模型自主二次决策:模型原始提问是查询北京、上海两地天气,人工把上海修改为西安后,模型检测到缺少上海的结果,会自主发起新一轮天气工具调用,流程再次触发中断;
结果差异化说明:不同大模型的自主决策逻辑存在差异,重复运行代码,模型是否二次调用工具、调用哪个城市,输出结果会有区别,属于正常现象。
补充说明: HumanInTheLoop 中间件通过可配置策略、灵活的决策类型和状态持久化,使 Agent 能够在关键操作上获得人工监督,既保证了自动化效率,又增加了安全性和可控性。结合流式处理,开发者可以构建出既流畅又可靠的交互体验。
流式传输子代理(需要掌握多 Agent 知识点)
当一个代理(supervisor)调用另一个代理(sub-agent)时,流式输出需要区分消息的来源,以便正确溯源。
实现原理:
- 在创建子代理时,使用
name参数为其命名。 - 在主代理流式调用时,设置
subgraphs=True。 - 在
messages模式的元数据中,通过lc_agent_name获取当前产生消息的代理名称。
代码示例:
from langchain.agents import create_agent
from langchain.chat_models import init_chat_model
from langchain.tools import tool
from langchain_core.messages import AIMessageChunk, AIMessage, ToolMessage
from langgraph.checkpoint.memory import InMemorySaver
import os
# ==================== 1. 定义天气工具 ====================
@tool
def get_weather(city: str) -> str:
"""
查询指定城市的天气信息
Args:
city: 城市名称,如"北京"、"上海"
Returns:
天气信息字符串
"""
weather_data = {
"北京": "晴天,温度25°C,湿度40%",
"上海": "多云,温度28°C,湿度65%",
"广州": "雷阵雨,温度30°C,湿度80%",
"深圳": "晴转多云,温度27°C,湿度70%",
"波士顿": "阴天,温度12°C,湿度75%",
"纽约": "晴朗,温度18°C,湿度55%",
"伦敦": "小雨,温度10°C,湿度85%"
}
import time
time.sleep(0.5) # 模拟网络延迟
return weather_data.get(city, f"未找到{city}的天气信息,请检查城市名称")
# ==================== 2. 创建子代理(天气代理) ====================
weather_agent = create_agent(
model=init_chat_model("gpt-4o-mini"), # 使用更经济的模型
tools=[get_weather],
name="weather_agent",
system_prompt="你是一个天气查询助手,负责查询指定城市的天气信息。",
)
# ==================== 3. 创建主代理 ====================
# 包装子代理为工具函数
def call_weather_agent(query: str) -> str:
"""调用天气代理查询天气信息"""
result = weather_agent.invoke({"messages": [{"role": "user", "content": query}]})
return result["messages"][-1].content
supervisor = create_agent(
model=init_chat_model("gpt-4o-mini"),
tools=[call_weather_agent],
name="supervisor",
system_prompt="你是一个主控代理,负责将用户的天气查询任务委派给天气代理。",
)
# ==================== 4. 流式传输主程序 ====================
def main():
print("===== 主代理开始执行 =====")
current_agent = None
config = {"configurable": {"thread_id": "1"}}
# 使用 subgraphs=True 来获取子代理的流式输出
for chunk in supervisor.stream(
{"messages": [{"role": "user", "content": "波士顿和纽约的天气如何?"}]},
config=config,
stream_mode=["messages", "updates"], # 同时获取消息和更新
subgraphs=True, # 重要:允许子图流式输出
version="v2",
):
# 处理消息流
if chunk["type"] == "messages":
token, meta = chunk["data"]
# 尝试获取代理名称
agent_name = meta.get("lc_agent_name")
# 如果代理名称变化,打印分隔符
if agent_name and agent_name != current_agent:
print(f"\n--- [{agent_name}] 正在输出: ---")
current_agent = agent_name
# 输出文本内容
if isinstance(token, AIMessageChunk) and token.text:
print(token.text, end="", flush=True)
# 处理更新(用于调试)
elif chunk["type"] == "updates":
for source, update in chunk["data"].items():
if source == "model":
# 模型输出完成
messages = update.get("messages", [])
if messages and isinstance(messages[-1], AIMessage):
print("\n[模型完成输出]")
elif source == "tools":
# 工具调用完成
messages = update.get("messages", [])
if messages and isinstance(messages[-1], ToolMessage):
print(f"\n[工具响应: {messages[-1].content[:50]}...]")
if __name__ == "__main__":
main()
关键点
subgraphs=True确保子代理的内部流式事件也能被捕获。- 元数据中的
lc_agent_name是区分代理来源的关键。 - 可在流中切换标签,使前端能明确区分不同代理的输出。
练习与思考
- 尝试在推理 Token 示例中,分别输出思考内容和最终文本,并观察它们出现的顺序。
- 编写同时处理生成文本增量
messages和完整调用信息updates,并在界面实时展示的示例。 - 设计一个需要人工审批的工具(如 “发送邮件”),实现中断、收集决策、恢复的完整流程。
更多推荐



所有评论(0)