Part 4用LangGraph的interrupt()实现条件性暂停,在周期初期和严重失衡时引入人工复核,防止自动化决策造成重大损失。
这是系列文章的第四部分。第三部分赋予了 agent 记忆能力,它现在可以连续运行数小时,一轮又一轮,不会忘记已经尝试过的方案。
这恰恰成了问题所在。一个真正严重的供需失衡可能需要激进的溢价系数,或者代价高昂的司机激励。当前 agent 会立即应用它决定的任何方案。在那之前没有人审核过。没有任何人会因为昂贵的后果而暂停来寻求第二意见。
第四部分引入了这个暂停机制。它使用 LangGraph 的动态 interrupt():
interrupt(payload)——但仅在确实需要暂停时才调用app.invoke(Command(resume=answer), config=thread) 从那里继续执行,answer 成为 interrupt() 的返回值,就在调用处两个门控暂停,而非一个:
早期——request_data_edit。门控条件为 cycle_number == 1。人类在 detect_imbalance 运行之前验证或更正一次原始起始快照。之后轮次使用 agent 自己携带的数值,而非新鲜未验证的读数——因此每轮都重新确认只会增加噪音。
晚期——request_approval。门控条件为 severity == "critical"。错误的自动决策在最严重的失衡中代价最高,所以这是人工审核值得其成本的地方。轻度、中度、平衡和盈余的轮次仍然完全自主运行,与第三部分完全相同。
一个值得指出的特性:没有向 state 添加新字段。这里的每个 interrupt 只读取第三部分已有的 state 字段,并向第三部分已有的字段写入——zone、recommended_policy、explanation。人工审核是一种新的控制流能力,而非新数据。
start_cycle ──▶ apply_scheduled_conditions ──▶ request_data_edit ──▶ detect_imbalance ──▶ classify_severity ──▶ set_candidates
(gate: cycle_number == 1) │
┌──────────┴──────────┐
▼ (balanced) ▼ (deficit/surplus)
trivial_do_nothing reconcile_inputs ← LLM #1
│ ▼
│ resolved_imbalance
│ ▼
│ choose_best_policy
│ ▼
│ generate_explanation ← LLM #2
└──────────┬──────────┘
▼
request_approval
(gate: severity == "critical")
▼
simulate_and_report
从 detect_imbalance 到 generate_explanation 的所有内容与第三部分完全相同——相同节点,直接复用。两个新节点夹在核心流程两端:request_data_edit 在轮次开始后立即执行,request_approval 在结果模拟之前执行。
两者都是纯 Python。在门控不适用的任何轮次中,两者都返回 {}——因此无门控轮次看起来与第三部分完全相同。
每个都是带顶部 if 守卫的普通函数。是否暂停不由 LangGraph 机制决定。节点自己决定,然后在需要时调用 interrupt(payload):
def request_data_edit(state: AgentState) -> AgentState:
if state["cycle_number"] != 1:
return {}
zone = state["zone"]
answer = interrupt({
"kind": "data_edit",
"question": (
f"Review the starting snapshot for {zone['zone_name']}. "
"Resume with {'corrections': {...}} to fix fields, or {'corrections': {}} to accept it."
),
"zone": zone,
})
corrections = answer.get("corrections", {})
if not corrections:
return {}
return {"zone": {**zone, **corrections}}
interrupt() 的返回值成为传递给 Command(resume=...) 的任意内容——这里是一个 corrections 字典。
一个真正的陷阱,在使用前值得了解:Command(resume={}) 在 Python 中是假值。LangGraph 将假值的 resume 视为没有给出任何答案。它只是重新触发同一个 interrupt,而非继续执行。始终用非空字典 resume——{"corrections": {}} 是你接受原快照的方式。
晚期暂停具有相同的形态,但有一个真正的决策需要做:
def request_approval(state: AgentState) -> AgentState:
if state["severity"] != "critical":
return {}
answer = interrupt({
"kind": "approval",
"question": (
f"Critical imbalance in {state['zone']['zone_name']} (cycle {state['cycle_number']}). "
"Approve, reject, or override the recommended policy."
),
"zone_name": state["zone"]["zone_name"],
"cycle": state["cycle_number"],
"imbalance_ratio": state["imbalance_ratio"],
"recommended_policy": state["recommended_policy"],
"explanation": state["explanation"],
"policy_evaluations": state["policy_evaluations"],
"candidate_policies": state["candidate_policies"],
})
action = (answer or {}).get("action", "approve")
if action == "reject":
return {
"recommended_policy": "do_nothing",
"explanation": "Rejected by human reviewer — falling back to do_nothing.",
}
if action == "override":
chosen = answer["policy"]
return {
"recommended_policy": chosen,
"explanation": f"Overridden by human reviewer: {chosen} chosen instead of "
f"{state['recommended_policy']}.",
}
return {} # approve — no change
三种可能的 resume 值:
{"action": "approve"}——保留推荐的策略{"action": "reject"}——回退到 do_nothing{"action": "override", "policy": "<name>"}——强制使用特定策略与第三部分相同的图,只是在这两个节点在正确的位置插入。它被包装在 build(checkpointer) 函数中,而非内联构建一次——下一节需要一个共享完全相同 checkpointer 的第二个图对象。结构从不改变。只是它是否是同一个 Python 对象有所不同:
def build(checkpointer):
llm = ChatOllama(model="qwen2.5:14b", temperature=0)
llm_with_reconcile_tool = llm.bind_tools([report_context_and_schedule])
def _reconcile(state):
return reconcile_inputs(state, llm_with_reconcile_tool)
def _explain(state):
return generate_explanation(state, llm)
g = StateGraph(AgentState)
g.add_node("start_cycle", start_cycle)
g.add_node("apply_scheduled_conditions", apply_scheduled_conditions)
g.add_node("request_data_edit", request_data_edit)
g.add_node("detect_imbalance", detect_imbalance)
g.add_node("classify_severity", classify_severity)
g.add_node("set_candidates", set_candidates)
g.add_node("trivial_do_nothing", trivial_do_nothing)
g.add_node("reconcile_inputs", _reconcile)
g.add_node("resolved_imbalance", resolved_imbalance)
g.add_node("choose_best_policy", choose_best_policy)
g.add_node("generate_explanation", _explain)
g.add_node("request_approval", request_approval)
g.add_node("simulate_and_report", simulate_and_report)
g.add_edge(START, "start_cycle")
g.add_edge("start_cycle", "apply_scheduled_conditions")
g.add_edge("apply_scheduled_conditions", "request_data_edit")
g.add_edge("request_data_edit", "detect_imbalance")
g.add_edge("detect_imbalance", "classify_severity")
g.add_edge("classify_severity", "set_candidates")
g.add_conditional_edges("set_candidates", route_llm_or_skip, {
"trivial_do_nothing": "trivial_do_nothing",
"reconcile_inputs": "reconcile_inputs",
})
g.add_edge("reconcile_inputs", "resolved_imbalance")
g.add_edge("resolved_imbalance", "choose_best_policy")
g.add_edge("choose_best_policy", "generate_explanation")
g.add_edge("generate_explanation", "request_approval")
g.add_edge("trivial_do_nothing", "request_approval")
g.add_edge("request_approval", "simulate_and_report")
g.add_edge("simulate_and_report", END)
return g.compile(checkpointer=checkpointer)
memory = MemorySaver()
app = build(memory)
第 1 轮总是在 request_data_edit 处暂停。Downtown Core 的原始快照在这里低估了司机数量——记录为 9,但实际是 15。
downtown = get_zone(zones, DEFAULT_ZONE_NAME, driver_count=9)
r1 = app.invoke(make_initial_state(downtown), config=thread_edit)
PAUSED — kind=data_edit
question: Review the starting snapshot for Downtown Core. Resume with {'corrections': {...}} to fix fields, or {'corrections': {}} to accept it.
zone as recorded: driver_count=9, rider_request_count=16
uncorrected ratio: 1.78
1.78 的比例显示为失衡。但这是一个错误的失衡——是数据问题,而非真实的供给问题。修正后继续执行:
MY_CORRECTION = {"driver_count": 15}
r2 = app.invoke(Command(resume={"corrections": MY_CORRECTION}), config=thread_edit)
(no interrupt — cycle ran straight through)
ratio now: 1.07 (balanced)
[Cycle 1] [Downtown Core] ratio=1.07 | severity=none | policy=do_nothing | wait 3.3min → 4.6min | resolved=N/A
一次修正,区域从失衡变为平衡。这就是这个暂停的真正意义——在错误读数驱动真实决策之前捕获它。
request_data_edit 门控条件为 cycle_number == 1,因此在同一线程的第 2 轮不会再触发。编辑门控每个线程只应用一次。
request_approval 仅在 severity == "critical" 时触发——比例已经糟糕到错误的自动决策代价高昂的程度。
到达暂停点,在一个有 4 名司机和 80 个乘客请求的区域:
PAUSED — kind=approval
question: Critical imbalance in Airport (cycle 1). Approve, reject, or override the recommended policy.
severity=critical | recommended=surge_pricing | candidates=['surge_pricing', 'driver_bonus', 'demand_redirect', 'do_nothing']
explanation: The surge pricing policy was selected for the Airport zone in Cycle 1 because it generated the highest profit of $142.93, even though it did not resolve the imbalance.
批准并继续执行——但这次使用不同的图对象:
MY_DECISION = {"action": "approve"}
app_resume = build(memory) # a different graph object, sharing only `memory`
r3 = app_resume.invoke(Command(resume=MY_DECISION), config=thread_approve)
final policy: surge_pricing
[Cycle 1] [Airport] ratio=20.0 | severity=critical | policy=surge_pricing | wait 5.6min → 13.3min | resolved=NO
那个 app_resume 细节值得暂停思考。它从相同的 build() 函数新鲜构建,它与原始 app 只共享一件事:MemorySaver。暂停状态不在 app Python 对象内部——它在 checkpointer 内部。任何以相同方式构建、共享那个 checkpointer 的图对象都可以恢复暂停。
request_data_edit 和 request_approval 都是有门控的。它们只在特定条件下暂停,resume 值会改变接下来发生的事情。
调试暂停在本质上有所不同:
def debug_checkpoint(state: AgentState) -> AgentState:
interrupt({
"kind": "debug",
"question": "Inspect state before the approval gate. Resume with anything to continue.",
"state_snapshot": {
"cycle_number": state["cycle_number"],
"zone": state["zone"],
"imbalance_ratio": state["imbalance_ratio"],
"severity": state["severity"],
"candidate_policies": state.get("candidate_policies"),
"policy_evaluations": state.get("policy_evaluations"),
"recommended_policy": state.get("recommended_policy"),
"explanation": state.get("explanation"),
},
})
return {}
它的存在是为了在审批门控采取行动之前检查完整的飞行中 state——而非做出或修改决策。它通过 build_graph(debug_mode=True) 可选接入,默认图中不包含它,这样它就不会在每次运行时都增加一个暂停。
决策核心仍然没有改变。choose_best_policy 正是第一部分的函数,调用方式完全相同。第四部分增加的是人类坐在循环内部的新方式,而非只是在事后读取报告——一个真正的暂停,可以从完全不同的 Python 对象恢复,有门控所以只在真正值得引起注意时才中断。
仍然缺失的部分:这里的每个区域仍然是独立评估的。自第一部分就提到的从邻近区域调入司机的策略从未真正构建。第五部分从这里继续。
本系列代码:github.com/ebiarian/zone-balancing-ridesharing-langgraph-agent
下一篇——第五部分:同时协调两个区域,使调入一个区域的司机成为一个真正同时考虑双方利益的决策。