Skip to content

第 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 持久化中断状态

🛠️ 动手实践

  1. 实现带流式输出的 RAG Agent(每段检索结果实时展示)
  2. 构建资金转账审批流程,模拟人工确认
  3. 实现中断恢复测试:保存状态 → 中断 → 恢复 → 验证