中断允许您在特定点暂停图执行并等待外部输入后再继续。这实现了人机协作模式,在此模式下您需要外部输入才能继续。当触发中断时,LangGraph 使用其 持久化 层保存图状态,并无限期等待直到您恢复执行。
中断通过调用 interrupt() 函数来实现,您可以在图的任意节点中调用它。该函数接受任何可序列化为 JSON 的值,并将其暴露给调用者。当您准备继续时,使用 Command重新调用图来恢复执行,这将成为 interrupt() 从节点内部调用的返回值。
与静态断点(在特定节点之前或之后暂停)不同,中断是 **动态的**:它们可以放置在代码的任何位置,并且可以根据应用逻辑有条件地触发。
- 检查点保存您的位置: 检查点写入器会精确保存图的状态,以便您稍后可以恢复,即使处于错误状态。
- - **
thread_id是指针:** 设置config={"configurable": {"thread_id": ...}}以指定检查点要加载的状态。 - - **中断负载通过以下方式暴露
stream.interrupts:** 当使用 事件流 (graph.stream_events(..., version="v3")时,您传递给interrupt()的值显示在stream.interrupts,并且stream.interruptedisTrue当运行暂停等待输入时。
您选择的 thread_id 实际上就是您的持久游标。重用它会从同一检查点恢复;使用新值会启动一个具有空状态的新线程。
使用以下方式暂停 interrupt
interrupt 函数会暂停图执行并向调用者返回值。当您调用 interrupt 时,LangGraph 会保存当前图状态并等待您恢复执行并提供输入。
要使用 interrupt,您需要: 1. A **检查点** 持久保存图状态(生产环境中使用持久化检查点) 2. A **线程 ID** 在配置中,以便运行时知道从哪个状态恢复 3. 在您想暂停的位置调用 interrupt() (负载必须可序列化为 JSON)
from langgraph.types import interrupt
def approval_node(state: State):
# Pause and ask for approval
approved = interrupt("Do you approve this action?")
# When you resume, Command(resume=...) returns that value here
return {"approved": approved}
当您调用 interrupt, 发生的事情如下:
- **图执行被挂起** 在恰好调用
interrupt的位置 - **状态被保存** 使用检查点器,以便稍后可以恢复执行,在生产环境中,这应该是一个持久化检查点器(例如,由数据库支持)
- **值被返回** 给调用方,
stream.interrupts当使用 事件流 (graph.stream_events(..., version="v3")时),或在__interrupt__使用默认invoke()API 时;它可以是任何 JSON 可序列化值(字符串、对象、数组等)
- **图无限期等待** 直到你使用响应恢复执行
- **响应在恢复时被传递回** 节点,成为
interrupt()调用的返回值
恢复会中断
在中断暂停执行后,你可以通过使用包含恢复值的 Command 再次调用图来恢复它。恢复值被传递回 interrupt 调用,使节点能够使用外部输入继续执行。
驱动可能中断的图的推荐方式是 事件流 —— 它通过 stream.interrupts 和 stream.interrupted暴露中断,并通过 stream.output.
关于恢复的关键要点:
- - 你必须使用相同的 **线程 ID** 来恢复,该 ID 是在发生中断时使用的
- - 传递给
Command(resume=...)的值成为interrupt调用的返回值 - - 节点从
interrupt被调用的节点开头重新启动,所以interrupt之前的任何代码都会再次运行 - - 你可以传递任何 JSON 可序列化值作为恢复值
常见模式
中断解锁的关键在于能够暂停执行并等待外部输入。这对于多种用例非常有用,包括:
- - 审批工作流:在执行关键操作(API 调用、数据库更改、金融交易)之前暂停
- - 处理多个中断:在单次调用中恢复多个中断时,将中断 ID 与恢复值配对
- - 审查和编辑:让人类在继续之前审查和修改 LLM 输出或工具调用
- - 中断工具调用:在执行工具调用之前暂停,以审查和编辑工具调用
- - 验证人工输入:在继续下一步之前暂停以验证人工输入
使用人在回路 (HITL) 中断进行流式传输
在构建具有人在回路工作流的交互式代理时,你可以使用 事件流式传输 在处理中断的同时并发消费消息块和状态快照。
在循环中使用 graph.stream_events(..., version="v3") 返回的类型化投影,直到运行完成: - 通过 stream.messages - 按 token 流式传输 AI 响应 stream.values - 通过 stream.interrupted 观察每步状态快照 stream.interrupts - 通过 stream_events 检测中断 Command(resume=...) 并从 stream.interrupted 读取其负载
- - **
stream.messages**通过调用message.text恢复执行stream.subgraphs[*].messages. - - **
stream.values**并重复直到 - - **
stream.interrupted/stream.interrupts**为 falsestream.interrupts - - **
Command(resume=...)**: 传递为下一个stream_events输入以恢复;循环直到运行完成而不中断
处理多个中断
当并行分支同时中断时(例如,扇出到多个节点,每个节点调用 interrupt()),您可能需要在一次调用中恢复多个中断。 使用单次调用恢复多个中断时,请将每个中断ID映射到其恢复值。 这确保了每个响应在运行时与正确的中断配对。
批准或拒绝
中断最常见的用途之一是在关键操作之前暂停并请求批准。例如,您可能希望让人类批准API调用、数据库更改或任何其他重要决策。
from typing import Literal
from langgraph.types import interrupt, Command
def approval_node(state: State) -> Command[Literal["proceed", "cancel"]]:
# Pause execution; payload shows up on stream.interrupts (with stream_events) or result["__interrupt__"] (with invoke)
is_approved = interrupt({
"question": "Do you want to proceed with this action?",
"details": state["action_details"]
})
# Route based on the response
if is_approved:
return Command(goto="proceed") # Runs after the resume payload is provided
else:
return Command(goto="cancel")
当您恢复图时,传递 True 以批准,或 False 以拒绝:
# To approve
graph.stream_events(Command(resume=True), config=config, version="v3").output
# To reject
graph.stream_events(Command(resume=False), config=config, version="v3").output
Full example
审阅和编辑状态
有时您希望让人类在继续之前审阅和编辑图状态的部分内容。这对于纠正LLM、添加缺失信息或进行调整非常有用。
from langgraph.types import interrupt
def review_node(state: State):
# Pause and show the current content for review (payload surfaces on stream.interrupts)
edited_content = interrupt({
"instruction": "Review and edit this content",
"content": state["generated_text"]
})
# Update the state with the edited version
return {"generated_text": edited_content}
恢复时,提供编辑的内容:
graph.stream_events(
Command(resume="The edited and improved text"), # Value becomes the return from interrupt()
config=config,
version="v3",
).output
Full example
工具中的中断
您也可以直接在工具函数中放置中断。这使得工具本身在每次被调用时暂停以等待批准,并允许在执行前对工具调用进行人工审阅和编辑。
首先,定义一个使用interrupt:
from langchain.tools import tool
from langgraph.types import interrupt
@tool
def send_email(to: str, subject: str, body: str):
"""Send an email to a recipient."""
# Pause before sending; payload surfaces on stream.interrupts when using event streaming
response = interrupt({
"action": "send_email",
"to": to,
"subject": subject,
"body": body,
"message": "Approve sending this email?"
})
if response.get("action") == "approve":
# Resume value can override inputs before executing
final_to = response.get("to", to)
final_subject = response.get("subject", subject)
final_body = response.get("body", body)
return f"Email sent to {final_to} with subject '{final_subject}'"
return "Email cancelled by user"
当您希望审批逻辑与工具本身绑定,使其可在图的不同部分重复使用时,这种方法特别有用。LLM可以自然地调用工具,而中断会在工具被触发时暂停执行,让您可以批准、编辑或取消该操作。
Full example
from typing import TypedDict, Annotated, Literal
from langchain.tools import tool
from langchain_anthropic import ChatAnthropic
from langgraph.checkpoint.sqlite import SqliteSaver
from langgraph.graph import StateGraph, START, END
from langgraph.types import Command, interrupt
from langchain.messages import AnyMessage, SystemMessage, ToolMessage
class AgentState(TypedDict):
messages: Annotated[list[AnyMessage], operator.add]
@tool
def send_email(to: str, subject: str, body: str):
"""Send an email to a recipient."""
# Pause before sending; payload surfaces on stream.interrupts when using event streaming
response = interrupt({
"action": "send_email",
"to": to,
"subject": subject,
"body": body,
"message": "Approve sending this email?",
})
if response.get("action") == "approve":
final_to = response.get("to", to)
final_subject = response.get("subject", subject)
final_body = response.get("body", body)
# Actually send the email (your implementation here)
print(f"[send_email] to={final_to} subject={final_subject} body={final_body}")
return f"Email sent to {final_to}"
return "Email cancelled by user"
model = ChatAnthropic(model="claude-sonnet-4-6").bind_tools([send_email])
tools_by_name = {"send_email": send_email}
def agent_node(state: AgentState):
# LLM may decide to call the tool; interrupt pauses before sending
result = model.invoke(state["messages"])
return {"messages": [result]}
def tool_node(state: AgentState):
"""Performs the tool call"""
result = []
for tool_call in state["messages"][-1].tool_calls:
tool = tools_by_name[tool_call["name"]]
observation = tool.invoke(tool_call["args"])
result.append(ToolMessage(content=observation, tool_call_id=tool_call["id"]))
return {"messages": result}
def should_continue(state: AgentState) -> Literal["tool_node", END]:
"""Decide if we should continue the loop or stop based upon whether the LLM made a tool call"""
messages = state["messages"]
last_message = messages[-1]
if last_message.tool_calls:
return "tool_node"
return END
builder = StateGraph(AgentState)
builder.add_node("agent", agent_node)
builder.add_node("tool_node", tool_node)
builder.add_edge(START, "agent")
builder.add_conditional_edges("agent", should_continue, ["tool_node", END]) # Routes to "tools" or END
builder.add_edge("tool_node", "agent") # Loop back after tools
checkpointer = SqliteSaver(
sqlite3.connect("tool-approval.db", check_same_thread=False)
)
graph = builder.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "email-workflow"}}
initial = graph.stream_events(
{
"messages": [
{"role": "user", "content": "Send an email to alice@example.com about the meeting"}
]
},
config=config,
version="v3",
)
initial.output # drive the stream to completion
print(initial.interrupts) # -> (Interrupt(value={'action': 'send_email', ...}),)
# Resume with approval and optionally edited arguments
resumed = graph.stream_events(
Command(resume={"action": "approve", "subject": "Updated subject"}),
config=config,
version="v3",
)
print(resumed.output["messages"][-1]) # -> Tool result returned by send_email
验证人工输入
有时您需要验证人工输入的值,若无效则重新提示。建议的方法是调用 interrupt() **每节点调用一次**,从节点返回并将错误信息存储在状态中,然后使用 **条件边** 循环回到节点,直到提供有效值。
正确的模式:
- 将重新提示的问题存储在状态中(例如。
pending_question). - 在节点中,调用
interrupt()**恰好一次**,传递状态中的当前问题。 - 如果答案无效,返回更新后的
pending_question以便下次调用时重新提示。 - 使用
add_conditional_edges路由回节点,直到收集到有效值。
每次恢复调用 get_age_node 恰好一次,运行 interrupt() 调用一次,然后退出。当答案无效时,条件边会循环返回,下次中断会用更新后的问题重新提示。每次恢复不会有代码运行超过一次。
Full example
中断规则
当你在节点内调用 interrupt 时,LangGraph 通过抛出一个特殊异常来暂停执行,该异常向运行时发出暂停信号。这个异常会向上传播通过调用栈,被运行时捕获,运行时通知图保存当前状态并等待外部输入。
当执行恢复时(在你提供请求的输入之后),运行时从头重新启动整个节点——它不会从调用 interrupt 的确切行继续。这意味着在 interrupt 之前运行的任何代码都会再次执行。正因为如此,在使用中断时需要遵循一些重要规则以确保它们按预期行为。
不要包装 interrupt calls in try/except
interrupt 通过抛出一个特殊异常来在调用点暂停执行。如果你用 try-catch 包装 interrupt call in a try/except block, you will catch this exception and the interrupt will not be passed back to the graph.
- * ✅ 将
interrupt调用与易出错的代码分开 - * ✅ Use specific exception types in try/except blocks
def node_a(state: State):
# ✅ Good: interrupting first, then handling
# error conditions separately
interrupt("What's your name?")
try:
fetch_data() # This can fail
except Exception as e:
print(e)
return state
def node_a(state: State):
# ✅ Good: catching specific exception types
# will not catch the interrupt exception
try:
name = interrupt("What's your name?")
fetch_data() # This can fail
except NetworkException as e:
print(e)
return state
- * 🔴 不要包装
interruptcalls in bare try/except blocks
def node_a(state: State):
# ❌ Bad: wrapping interrupt in bare try/except
# will catch the interrupt exception
try:
interrupt("What's your name?")
except Exception as e:
print(e)
return state
不要重新排序 interrupt 节点内的调用
在单个节点中使用多个中断很常见,但如果处理不当可能会导致意外行为。
当节点包含多个 interrupt 调用时,LangGraph 会为执行该节点的任务维护一个特定的恢复值列表。每当执行恢复时,它会从节点的开头开始。对于遇到的每个 interrupt,LangGraph 会检查任务的恢复列表中是否存在匹配的值。匹配是 **严格基于索引的**,因此节点内 interrupt 调用的顺序很重要。
- * ✅ 保持
interrupt调用在节点执行过程中保持一致
def node_a(state: State):
# ✅ Good: interrupt calls happen in the same order every time
name = interrupt("What's your name?")
age = interrupt("What's your age?")
city = interrupt("What's your city?")
return {
"name": name,
"age": age,
"city": city
}
- * 🔴 不要有条件地跳过
interrupt在节点内的调用 - * 🔴 不要循环
interrupt使用在执行过程中不确定的逻辑进行调用,包括while True验证循环。请改用条件边(请参阅 验证人工输入)
def node_a(state: State):
# ❌ Bad: conditionally skipping interrupts changes the order
name = interrupt("What's your name?")
# On first run, this might skip the interrupt
# On resume, it might not skip it - causing index mismatch
if state.get("needs_age"):
age = interrupt("What's your age?")
city = interrupt("What's your city?")
return {"name": name, "city": city}
def node_a(state: State):
# ❌ Bad: looping based on non-deterministic data
# The number of interrupts changes between executions
results = []
for item in state.get("dynamic_list", []): # List might change between runs
result = interrupt(f"Approve {item}?")
results.append(result)
return {"results": results}
不要在 interrupt 调用中返回复杂值
根据所使用的检查点不同,复杂值可能无法序列化(例如,您无法序列化函数)。为了使图能够适应任何部署环境,最佳实践是仅使用可以合理序列化的值。
- * ✅ 传递简单的、JSON 可序列化的类型到
interrupt - * ✅ Pass dictionaries/objects with simple values
def node_a(state: State):
# ✅ Good: passing simple types that are serializable
name = interrupt("What's your name?")
count = interrupt(42)
approved = interrupt(True)
return {"name": name, "count": count, "approved": approved}
def node_a(state: State):
# ✅ Good: passing dictionaries with simple values
response = interrupt({
"question": "Enter user details",
"fields": ["name", "email", "age"],
"current_values": state.get("user", {})
})
return {"user": response}
- * 🔴 不要将函数、类实例或其他复杂对象传递给
interrupt
def validate_input(value):
return len(value) > 0
def node_a(state: State):
# ❌ Bad: passing a function to interrupt
# The function cannot be serialized
response = interrupt({
"question": "What's your name?",
"validator": validate_input # This will fail
})
return {"name": response}
class DataProcessor:
def __init__(self, config):
self.config = config
def node_a(state: State):
processor = DataProcessor({"mode": "strict"})
# ❌ Bad: passing a class instance to interrupt
# The instance cannot be serialized
response = interrupt({
"question": "Enter data to process",
"processor": processor # This will fail
})
return {"result": response}
在 interrupt 之前调用的副作用必须是幂等的
由于 interrupt 的工作方式是通过重新运行调用它们的节点,因此在 interrupt 之前调用的副作用(理想情况下)应该是幂等的。从上下文来看,幂等性意味着同一操作可以多次应用,而不会改变初次执行之后的结果。
例如,您可能在节点内部有一个 API 调用来更新记录。如果 interrupt 在该调用之后被调用,那么当节点恢复时它会被重新运行多次,可能会覆盖初始更新或创建重复记录。
def node_a(state: State):
# ✅ Good: using upsert operation which is idempotent
# Running this multiple times will have the same result
db.upsert_user(
user_id=state["user_id"],
status="pending_approval"
)
approved = interrupt("Approve this change?")
return {"approved": approved}
def node_a(state: State):
# ✅ Good: placing side effect after the interrupt
# This ensures it only runs once after approval is received
approved = interrupt("Approve this change?")
if approved:
db.create_audit_log(
user_id=state["user_id"],
action="approved"
)
return {"approved": approved}
def approval_node(state: State):
# ✅ Good: only handling the interrupt in this node
approved = interrupt("Approve this change?")
return {"approved": approved}
def notification_node(state: State):
# ✅ Good: side effect happens in a separate node
# This runs after approval, so it only executes once
if (state.approved):
send_notification(
user_id=state["user_id"],
status="approved"
)
return state
- * 🔴 不要在
interrupt - * 🔴 不要在创建新记录前不检查其是否存在
def node_a(state: State):
# ❌ Bad: creating a new record before interrupt
# This will create duplicate records on each resume
audit_id = db.create_audit_log({
"user_id": state["user_id"],
"action": "pending_approval",
"timestamp": datetime.now()
})
approved = interrupt("Approve this change?")
return {"approved": approved, "audit_id": audit_id}
def node_a(state: State):
# ❌ Bad: appending to a list before interrupt
# This will add duplicate entries on each resume
db.append_to_history(state["user_id"], "approval_requested")
approved = interrupt("Approve this change?")
return {"approved": approved}
将 interrupt 与作为函数调用的子图结合使用
在节点内调用子图时,父图将从 **节点的开头** (子图被调用的位置)恢复执行,并且 interrupt 被触发时也会发生类似情况 **子图** 也将从节点的开头恢复执行interrupt被调用的位置)
def node_in_parent_graph(state: State):
some_code() # <-- This will re-execute when resumed
# Invoke a subgraph as a function.
# The subgraph contains an `interrupt` call.
subgraph_result = subgraph.invoke(some_input)
# ...
def node_in_subgraph(state: State):
some_other_code() # <-- This will also re-execute when resumed
result = interrupt("What's your name?")
# ...
使用 interrupt 进行调试
要调试和测试图,您可以使用静态 interrupt 作为断点来逐步执行图,每次一个节点。静态 interrupt 在定义的点触发,要么在节点执行之前,要么在节点执行之后。您可以通过指定 interrupt_before 和 interrupt_after 来设置这些参数
At compile time
graph = builder.compile(
interrupt_before=["node_a"], # [!code highlight]
interrupt_after=["node_b", "node_c"], # [!code highlight]
checkpointer=checkpointer,
)
# Pass a thread ID to the graph
config = {
"configurable": {
"thread_id": "some_thread"
}
}
# Run the graph until the breakpoint
graph.invoke(inputs, config=config) # [!code highlight]
# Resume the graph
graph.invoke(None, config=config) # [!code highlight]
- 断点在
compiletime. interrupt_before指定执行应在节点执行之前暂停的节点。interrupt_after指定执行应在节点执行之后暂停的节点。- 需要检查点器来启用断点。
- 图运行直到命中第一个断点。
- 通过传入
None作为输入来恢复图。这将运行图直到命中下一个断点。
At run time
config = {
"configurable": {
"thread_id": "some_thread"
}
}
# Run the graph until the breakpoint
graph.invoke(
inputs,
interrupt_before=["node_a"], # [!code highlight]
interrupt_after=["node_b", "node_c"], # [!code highlight]
config=config,
)
# Resume the graph
graph.invoke(None, config=config) # [!code highlight]
graph.invoke使用interrupt_before和interrupt_after参数调用。这是一种运行时配置,可以在每次调用时更改。interrupt_before指定执行应在节点执行之前暂停的节点。interrupt_after指定执行应在节点执行之后暂停的节点。- 图运行直到命中第一个断点。
- 通过传入
None作为输入来恢复图。这将运行图直到命中下一个断点。
使用 LangSmith Studio
您可以使用 LangSmith Studio 在 UI 中运行图之前设置静态中断。您还可以使用 UI 检查执行过程中任何点的图状态。
!图片