invoke和stream概述
invoke和stream都是编译后 graph 的执行入口,用来启动整个图运行;底层真正干活的是同一套 Pregel BSP 引擎。
invoke:一次性同步阻塞调用,全部 Superstep 跑完,返回最终完整 state。stream:返回迭代器,每完成一个 Superstep,立刻产出这一轮的增量片段,可以拿到中间过程。
二者底层跑的是同一套图逻辑,执行结果完全一致,只是返回数据时机不一样。
graph.invoke(input)
函数签名
final_state = graph.invoke( input={"messages": [("user", "帮我执行压测,分析jmeter报告")]}, config={"configurable": {"thread_id":"test‑001"}} )执行行为
- 同步阻塞,主线程卡住,整个图全部执行完毕才返回;所有 Superstep 完整跑完(循环、重试全部跑完)。
- 返回:完整最终的 AgentState 全部状态字典。包含所有 messages、trajectory、中间变量。
- thread_id 配合
Checkpointer,会自动保存每一步检查点快照。
优点
- 代码极简,非常适合 Eval‑Harness 自动化评测 Eval‑Harness 跑自动化用例,不需要中间输出,只需要最终结果,直接拿
final_state做指标判断(ToolCorrectness、结果正确性)。
# Eval‑Harness典型写法 result_state = graph.invoke(golden_case_input, config=config) actual_output = result_state["messages"][-1].content # deepeval metric直接评测 actual_output- 不需要处理迭代循环,不需要解析中间片段,写自动化测试脚本非常清爽。
缺点
- 看不到中间执行轨迹:运行过程中看不到:哪个节点跑了、工具调用、中间报错,全部等跑完才拿到。
- 如果 Agent 循环很多、LLM 调用耗时久,会长时间阻塞,控制台没有任何输出,看起来像卡死。
- 如果图内部设置
interrupt()人工暂停:invoke遇到中断会直接抛出异常,invoke 不适合处理人机暂停场景。
关键点:invoke 遇到interrupt()会报错,因为 invoke 设计是 “一口气跑完”,遇到暂停无法继续。处理人工介入只能用stream。
适用场景
- Eval‑Harness 自动化回归测试(你的场景)
- 后台服务,只需要最终输出,不展示中间过程
- 脚本批量跑用例
graph.stream(input)
函数签名
stream_output = graph.stream( input={"messages": [("user", "帮我执行压测,分析jmeter报告")]}, config={"configurable": {"thread_id":"test‑001"}} ) for chunk in stream_output: print(chunk) # 每一个Superstep结束,产出一轮chunk增量执行行为
- 返回一个 Python 迭代器,不会阻塞等待全部完成。
- 每完成 1 个 Superstep(一轮 BSP 超级步、屏障同步完成),就产出一个 chunk 片段。
- chunk 不是完整 state!chunk 只是本轮 Superstep 产生的增量更新,只包含本轮被修改过的字段。
- for 循环每一轮拿到的
chunk格式示例:
# chunk示例,key=节点名称,value=该节点返回的增量字典 { "supervisor": { "next_agent": "lighthouse_agent", "messages": [AIMessage(...)] } }chunk 只有本轮变更字段,不会携带全部历史 state。如果你想要完整 state,需要借助 checkpointer:graph.get_state(config)读取完整快照。
优点
- 实时观测中间 Trajectory 轨迹 每跑完一个节点(Superstep)立刻拿到输出,可以打印:哪个节点执行、工具调用参数、中间报错,调试 Agent 必备。
- 原生支持识别
interrupt()人工暂停 当 Superstep 边界触发中断,stream 迭代器产出__interrupt__标记,迭代停止;外部系统可以展示 UI 给操作人员,调用graph.resume()继续执行。 - 可以做前端流式 UI:一边跑 Agent,一边把中间步骤展示给用户。
缺点
- 拿到的是增量片段,不是完整 state,业务代码需要自己拼装完整上下文;
- 写 Eval‑Harness 自动化测试会多写一层 for 循环,代码更啰嗦;
- chunk 结构会随节点变化,解析要做判断。
获取完整状态小技巧,stream 循环内随时读取完整快照:
for chunk in graph.stream(input, config): full_snapshot = graph.get_state(config) # 读取checkpointer里完整state print(full_snapshot.values["messages"])适用场景
- Agent 调试,观察每一步 Superstep 执行轨迹;
- 前端流式交互,展示 Agent 思考、工具调用过程;
- Human‑in‑the‑loop,需要人工暂停、确认再继续;
- 需要采集每一步 trajectory 日志。
关键对比表格
| 项目 | graph.invoke() | graph.stream() |
|---|---|---|
| 返回值 | 完整最终 State 字典 | 迭代器,每 Superstep 返回增量 chunk |
| 阻塞行为 | 完全阻塞,全部跑完返回 | 迭代产出,不会阻塞等待全部完成 |
| 拿到的数据 | 全部字段 | 仅本轮修改的增量字段 |
| 处理 interrupt 中断 | 遇到中断直接抛异常 | 识别中断标记,迭代停止 |
| 获取中间 Trajectory | 拿不到中间过程 | 每一轮 Superstep 拿到中间输出 |
| Eval‑Harness 测试 | 首选,代码简洁 | 可以用,但需要额外 get_state 拿完整状态 |
| 调试 Agent | 不推荐 | 首选,看每一步节点输出 |
高频踩坑点
坑 1:stream 拿到 chunk 以为是完整 state
# 错误写法:chunk只是增量,没有全部字段 for chunk in graph.stream(input, config): print(chunk["messages"]) # 有可能KeyError,本轮没有更新messages就不存在 # 正确:想要完整状态,调用 graph.get_state(config).values坑 2:invoke 为什么不能处理 interrupt
interrupt()生效点在Superstep 屏障同步边界。invoke 希望一口气跑完所有 Superstep;一旦遇到中断,没有对外暴露暂停接口,直接抛出错误。想要人机交互,必须用 stream。
坑 3:同一个 thread_id,invoke 和 stream 可以混用
checkpointer 保存快照,invoke 跑完之后,可以用同一个 thread_id 调用 stream 继续执行;反之也可以。
Eval‑Harness 测试实践建议
自动化回归测试用例:优先使用 invoke
# golden case自动化测试,不需要中间步骤 final_state = graph.invoke(test_input, config=thread_config) # final_state拿到完整state,提取输出、工具调用记录,交给deepeval指标评测本地调试 Agent,看 trajectory 轨迹:使用 stream + get_state ()
for chunk in graph.stream(test_input, config=thread_config): snapshot = graph.get_state(config=thread_config).values print("==== Superstep完成 ====") print("当前轨迹messages:", snapshot["messages"])时序小例子(Supervisor 多 Agent)
流程:START → supervisor → lighthouse_agent → supervisor → END一共 3 个 Superstep
- Superstep1:执行 supervisor;stream 产出 chunk:
{"supervisor": {...}} - Superstep2:执行 lighthouse_agent;stream 产出 chunk:
{"lighthouse_agent": {...}} - Superstep3:再次执行 supervisor;stream 产出 chunk:
{"supervisor": {...}}
- stream:循环会收到 3 次 chunk;
- invoke:全部 3 轮跑完,一次性返回合并后的完整 state。
实例Demo
下面是完整可运行最小 Demo,实现一个简单的反思循环 Agent,同时演示graph.invoke()和graph.stream(),打印对比输出。
依赖:pip install langgraph,不需要大模型,全部模拟逻辑,直接跑就能看到效果。
# -*- coding:utf-8 -*- from typing import TypedDict, Annotated import operator from langgraph.graph import StateGraph, START, END from langgraph.checkpoint.memory import MemorySaver # 1. 定义状态 class State(TypedDict): messages: Annotated[list, operator.add] count: int # 2. 定义节点(模拟业务逻辑,不调用LLM) def think_node(state: State): """思考节点:计数+1,模拟agent思考""" cnt = state["count"] new_msg = f"第{cnt+1}次思考完成" return { "messages": [new_msg], "count": cnt + 1 } def judge_node(state: State): """判断节点:最多循环3次就结束""" cnt = state["count"] if cnt >= 3: return {"next": "end"} else: return {"next": "loop"} # 3. 构建图 builder = StateGraph(State) builder.add_node("think", think_node) builder.add_node("judge", judge_node) builder.add_edge(START, "think") builder.add_edge("think", "judge") # 条件边:循环回退到 think builder.add_conditional_edges( "judge", lambda s: s["next"], { "loop": "think", "end": END } ) checkpointer = MemorySaver() graph = builder.compile(checkpointer=checkpointer) config = {"configurable": {"thread_id": "demo‑001"}} init_input = {"messages": [], "count": 0} if __name__ == "__main__": print("========== 【1】演示 graph.stream() 每一个Superstep返回增量chunk ==========\n") stream_iter = graph.stream(init_input, config=config) for chunk in stream_iter: print(f" stream收到chunk(本轮增量): {chunk}") # 读取checkpoint里面的完整快照状态 full_snap = graph.get_state(config).values print(f" 当前完整state: count={full_snap['count']}, messages={full_snap['messages']}\n") print("\n========== 重置thread_id,执行 graph.invoke() 阻塞等待全部完成,直接拿最终完整state ==========\n") config_invoke = {"configurable": {"thread_id": "demo‑002"}} final_state = graph.invoke(init_input, config=config_invoke) print(f"invoke返回最终完整state:") print(f"count = {final_state['count']}") print(f"messages = {final_state['messages']}")输出样例(控制台打印)
========== 【1】演示 graph.stream() 每一个Superstep返回增量chunk ========== stream收到chunk(本轮增量): {'think': {'messages': ['第1次思考完成'], 'count': 1}} 当前完整state: count=1, messages=['第1次思考完成'] stream收到chunk(本轮增量): {'judge': {'next': 'loop'}} 当前完整state: count=1, messages=['第1次思考完成'] stream收到chunk(本轮增量): {'think': {'messages': ['第2次思考完成'], 'count': 2}} 当前完整state: count=2, messages=['第1次思考完成', '第2次思考完成'] stream收到chunk(本轮增量): {'judge': {'next': 'loop'}} 当前完整state: count=2, messages=['第1次思考完成', '第2次思考完成'] stream收到chunk(本轮增量): {'think': {'messages': ['第3次思考完成'], 'count': 3}} 当前完整state: count=3, messages=['第1次思考完成', '第2次思考完成', '第3次思考完成'] stream收到chunk(本轮增量): {'judge': {'next': 'end'}} 当前完整state: count=3, messages=['第1次思考完成', '第2次思考完成', '第3次思考完成'] ========== 重置thread_id,执行 graph.invoke() 阻塞等待全部完成,直接拿最终完整state ========== invoke返回最终完整state: count = 3 messages = ['第1次思考完成', '第2次思考完成', '第3次思考完成']重点观察现象
stream
- 每一轮 Superstep(节点执行完毕 + 屏障同步)产出一个
chunk chunk只是本轮节点返回的增量字典,不是全部 state- 如果想要完整数据,需要调用
graph.get_state(config).values读取检查点快照 - 循环过程中间每一步都可以打印,适合调试、采集 trajectory
invoke
- 程序卡住阻塞,所有 Superstep 全部执行完毕才返回
- 返回直接就是完整合并后的最终 state,看不到中间每一轮增量
- 适合 Eval‑Harness 自动化测试,代码干净,直接拿结果做指标判断
一句话总结
invoke = 等全部戏演完,一次性拿到完整剧本; stream = 每演完一幕,就把这一幕的剧本片段递给你。