人机协作(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 中断该工具,允许全部三种决策(approveeditreject "工具1": True
False 不中断,自动批准(工具调用直接执行) "工具2": False
InterruptOnConfig 对象 精细控制:可指定允许的决策列表和自定义描述 见下方
InterruptOnConfig 对象的属性:
  • allowed_decisions列表,可选值 "approve""edit""reject"。决定人工可用的操作类型。

  • description字符串或可调用函数,用于覆盖该工具的中断提示消息。若未指定,则使用全局 description_prefix 拼接而成。

我们以三类业务工具举例,讲解interrupt_on字典的三种配置写法,分别对应不同业务风险等级:

示例工具列表

  1. write_file:写文件,会修改磁盘存储,属于高风险操作;
  2. execute_sql:执行 SQL 语句,可能删改数据库数据,中等风险;
  3. 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_ondescription_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 自动放行;对写操作建议至少配置 approvereject

代码中 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:该阶段更新后的完整消息状态,包含所有对话、工具返回内容。

执行打印后能清晰看到完整执行流程:

  1. 用户提问「北京的天气如何」,首先进入 model 节点,大模型生成调用天气工具的指令,此时状态里仅存在工具调用 AI 消息;

  2. 接着进入 tools 工具节点,执行天气工具,返回「北京阳光明媚」的结果;

  3. 最后再次进入 model 节点,大模型整合工具结果,生成最终回复文本。 这就是 updates 模式的完整作用:流式同步每一步执行后的更新状态。

==================================================
用户提问: 北京的天气如何
==================================================

开始流式输出...

【v2 流式输出】
----------------------------------------
分片类型: updates
  节点: model
  内容: 好的,我来查一下北京的天气情况!...
----------------------------------------
分片类型: updates
  节点: tools
  内容: 晴天,温度25°C,湿度40%...
----------------------------------------
分片类型: updates
  节点: model
  内容: 北京今天天气不错哦!☀️

- **天气状况**:晴天
- **温度**:25°C,体感舒适
- **湿度**:40%,空气比较干爽

非常适合出门活动!建议穿短袖或薄长袖,记得做好防晒哦~🌤️ 如果想了解其他城市的天气,随时问我!...
----------------------------------------

注意:将 version="v2"(需要 LangGraph >= 1.1.)传递给 stream()astream() 以获取统一的输出格式。每个数据块都是一个具有 typensdata 键的 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(对应负载)的字典。 示例:同时使用 updatescustom 模式

模式实践

接下来我们讲解流式能力几个高频落地实践场景,在掌握流模式基础用法后,我们梳理它在业务里最常用的落地场景。第一个核心场景就是流式传输模型推理、思考 Token

流式传输推理 / 思考 Token

某些模型(如 Claude、OpenAI 系列)在生成最终答案前会进行内部推理。将这些推理内容实时输出,可以增强应用的透明度与可解释性。

我们可以通过配置开启推理输出,把模型幕后的思考步骤实时展示给用户。 举个直观例子:我们向模型提问 “你好”,交互界面里会先出现一段 “正在思考” 的文字,这段内容就是大模型原生输出的推理过程,我们可以通过代码配置把这段思考内容完整捕获并流式打印。

实现原理:LangChain 将不同提供商的推理内容(Anthropic 的thinking块、OpenAI 的reasoning摘要等)统一规范为 content_blockstype: "reasoning" 的块。

代码示例:

以 OpenAI 系列模型为例,想要开启推理输出,初始化 ChatOpenAI 时需要新增 reasoning 配置项,内部包含两个核心参数:

  1. effort(推理程度):可选 lowmediumhigh 三档,控制模型思考的深度;
  2. summary(推理摘要):支持 detailed(详细完整推理)、auto(自动精简)、None(不生成摘要)三种模式。

完成模型推理开关配置后,我们结合 Agent 的流式调用,就能实现推理过程逐 Token 实时流式输出的效果。


以 DeepSeek 系列模型为例,想要开启推理输出,初始化 ChatOpenAI 时需要新增 extra_body 配置项,内部包含推理参数:

核心参数

reasoning_effort(推理程度):可选 lowmediumhigh 三档,控制模型思考的深度;

⚠️ 注意: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"]内的所有节点更新,仅处理modeltools两类核心节点:

  • 如果更新消息是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流模式后,我们能完整流式还原工具调用全流程:

  1. 实时分片流式输出工具调用的 JSON 参数;

  2. 模型节点执行完毕后,输出完整解析后的工具调用记录;

  3. 工具执行完成后,输出工具返回的完整结果;

  4. 最终流式输出大模型整合后的回答文本。

关键点

  • 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)时,流式输出需要区分消息的来源,以便正确溯源。

实现原理

  1. 在创建子代理时,使用 name 参数为其命名。
  2. 在主代理流式调用时,设置 subgraphs=True
  3. 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 是区分代理来源的关键。
  • 可在流中切换标签,使前端能明确区分不同代理的输出。

练习与思考

  1. 尝试在推理 Token 示例中,分别输出思考内容和最终文本,并观察它们出现的顺序。
  2. 编写同时处理生成文本增量 messages 和完整调用信息 updates,并在界面实时展示的示例。
  3. 设计一个需要人工审批的工具(如 “发送邮件”),实现中断、收集决策、恢复的完整流程。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐