第 8 章 · 流式输出与中断恢复
本章目标:
- 掌握 stream() 与 invoke() 的区别
- 理解 StreamMode 的不同模式
- 实现中断恢复(interrupt + resume)
- 构建带审批流程的 Agent
8.1 stream() 基础
python
from langgraph.graph import StateGraph, MessagesState, START, END
def llm_node(state: MessagesState) -> dict:
return {"messages": [{"role": "assistant", "content": "Hello!"}]}
graph = StateGraph(MessagesState)
graph.add_node("llm", llm_node)
graph.add_edge(START, "llm")
graph.add_edge("llm", END)
compiled = graph.compile()
# 流式执行
for event in compiled.stream({"messages": [{"role": "user", "content": "hi"}]}):
print(event)
# 输出示例: {'llm': {'messages': [...]}}8.2 StreamMode 模式
python
# values:完整状态(默认)
for event in compiled.stream(input, stream_mode="values"):
print(event["messages"][-1]["content"])
# messages:只输出消息 delta
for event in compiled.stream(input, stream_mode="messages"):
msg, metadata = event
print(msg.content, end="", flush=True)
# events:输出事件详情
for event in compiled.stream(input, stream_mode="events"):
print(event)8.3 中断机制
python
from langgraph.types import interrupt
def approval_node(state: MessagesState) -> dict:
# 中断等待审批
decision = interrupt({
"action": "transfer_funds",
"amount": 10000,
"question": "确认转账 10000 元?"
})
if decision.get("approved"):
return {"messages": [{"role": "assistant", "content": "转账成功"}]}
return {"messages": [{"role": "assistant", "content": "转账已取消"}]}
# 配置中断点
graph.add_node("approval", approval_node)
graph.add_edge("approval", END)8.4 中断恢复
python
# 保存 checkpoint_id
result = compiled.invoke(input, config={"configurable": {"thread_id": "t1"}})
checkpoint_id = result["__checkpoints__"]["t1"]["latest"]
# 恢复执行
from langgraph.types import Command
resumed = compiled.invoke(
Command(resume={"approved": True}),
config={"configurable": {"thread_id": "t1", "checkpoint_id": checkpoint_id}}
)8.5 完整审批流程示例
python
from langgraph.graph import StateGraph, MessagesState, START, END
from langgraph.types import interrupt
from typing import TypedDict
class TransferState(TypedDict):
messages: list
amount: int
recipient: str
approved: bool
def initiate_node(state: TransferState) -> dict:
return {
"messages": [{
"role": "assistant",
"content": f"准备转账 {state['amount']} 元到 {state['recipient']}"
}]
}
def approval_node(state: TransferState) -> dict:
decision = interrupt({
"type": "fund_transfer",
"amount": state["amount"],
"recipient": state["recipient"]
})
return {"approved": decision.get("approved", False)}
def execute_node(state: TransferState) -> dict:
if state["approved"]:
return {"messages": [{"role": "assistant", "content": "转账成功"}]}
return {"messages": [{"role": "assistant", "content": "转账已取消"}]}
graph = StateGraph(TransferState)
graph.add_node("initiate", initiate_node)
graph.add_node("approval", approval_node)
graph.add_node("execute", execute_node)
graph.add_edge(START, "initiate")
graph.add_edge("initiate", "approval")
graph.add_edge("approval", "execute")
graph.add_edge("execute", END)
compiled = graph.compile()本章小结
stream()支持实时输出,适合前端展示interrupt()实现人机协作审批- 恢复执行使用
Command(resume=...) - 生产环境需配合 Checkpointer 持久化中断状态
🛠️ 动手实践
- 实现带流式输出的 RAG Agent(每段检索结果实时展示)
- 构建资金转账审批流程,模拟人工确认
- 实现中断恢复测试:保存状态 → 中断 → 恢复 → 验证