怎么写?让 Agent 先规划,再执行,再补缺口
文章摘要:ReAct 擅长根据局部反馈选择下一步,却容易在长任务里漏项、重复搜索和过早总结。本文只讲 Plan-and-Execute:如何把计划变成结构化状态,并让预算、执行、审查、并行和重规划真正进入运行时。
所属系列:ReAct 与 Agent Loop 工程(3/6)
上一篇的 Harness 管一次 Agent 运行内部的预算、重试、审批和副作用。本篇再向外一层:Plan-and-Execute 管多个 worker 调用之间的全局覆盖、依赖关系和重规划。两者可以叠加,但不是同一层能力。
用户提出:“分析 5 家竞品最近两个版本,对比定价、核心功能和用户反馈,最后给出进入策略。”
如果直接把这句话扔给一个 ReAct Agent,它很可能搜到哪写到哪。最后的文字或许流畅,但你很难回答:5 家是否全部覆盖?结论是否有来源?漏项以后由谁补?
计划必须是数据,不是作文
一个可执行计划至少要包含步骤 ID、验收条件、依赖、工具预算和超时:
from pydantic import BaseModel, Field
class Step(BaseModel):
id: str
task: str
success_criteria: list[str]
depends_on: list[str] = Field(default_factory=list)
tool_budget: int = Field(default=4, ge=1, le=10)
timeout_seconds: int = Field(default=90, ge=10, le=600)
class Plan(BaseModel):
steps: list[Step] = Field(min_length=1, max_length=20)
class Evidence(BaseModel):
step_id: str
summary: str
source_ids: list[str] = Field(default_factory=list)
passed: bool
missing: list[str] = Field(default_factory=list)
error_type: str | None = None
class Review(BaseModel):
approved: bool
gaps: list[str] = Field(default_factory=list)
tool_budget 如果只出现在 schema 里而不进入执行器,就是假护栏。depends_on 如果不参与调度,也只是计划书里的装饰。
Plan-and-Execute 的控制流
目标 -> 结构化计划 -> 执行可运行步骤 -> 汇总证据 -> 审查
^ |
| | 不通过
+------ 只补缺口的重规划 -+
审查通过,或预算不再允许重规划 -> 生成最终交付物
执行节点内部仍然可以是 ReAct Agent。Plan-and-Execute 管全局覆盖,ReAct 负责一个步骤内部根据工具反馈行动。
一个会执行预算的 StateGraph
下面示例是串行基线版。重点不是业务搜索工具,而是四件事:每一步动态创建带工具上限的 worker;单步有墙钟超时;失败被转成结构化证据;重规划上限来自任务预算,而不是散落在路由函数里的魔法数字。
import asyncio
import operator
import os
from typing import Annotated, Literal, TypedDict
from langchain.agents import create_agent
from langchain.agents.middleware import ToolCallLimitMiddleware, ToolRetryMiddleware
from langchain.chat_models import init_chat_model
from langchain.tools import tool
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import END, START, StateGraph
from pydantic import ValidationError
MODEL = os.getenv("AGENT_MODEL")
if not MODEL:
raise RuntimeError("请把 AGENT_MODEL 设置为你实际有权限使用的模型标识")
@tool
async def search_public_sources(query: str) -> list[dict]:
"""搜索公开资料,返回可追溯的来源 ID。"""
return [{
"source_id": "demo-source-001",
"title": "产品更新记录",
"snippet": f"与查询“{query}”相关的演示资料",
}]
planner = init_chat_model(MODEL).with_structured_output(Plan)
reviewer = init_chat_model(MODEL).with_structured_output(Review)
writer = init_chat_model(MODEL)
def make_worker(tool_budget: int):
return create_agent(
model=MODEL,
tools=[search_public_sources],
response_format=Evidence,
system_prompt=(
"一次只执行一个计划步骤。逐条核对 success_criteria;"
"没有来源 ID 的事实结论不能标记为 passed。"
),
middleware=[
ToolCallLimitMiddleware(
run_limit=tool_budget,
exit_behavior="error",
),
ToolRetryMiddleware(max_retries=2),
],
)
class TaskState(TypedDict):
objective: str
plan: list[Step]
cursor: int
evidence: Annotated[list[Evidence], operator.add]
review: Review | None
replan_count: int
max_replans: int
final_answer: str
def failed_evidence(step: Step, error_type: str, detail: str) -> Evidence:
return Evidence(
step_id=step.id,
summary="该步骤执行失败,已交给审查器决定是否重规划。",
passed=False,
missing=[detail[:500]],
error_type=error_type,
)
async def plan_node(state: TaskState) -> dict:
try:
plan = await planner.ainvoke(
f"目标:{state['objective']}\n"
"生成结构化计划,每步必须包含可检查的验收条件。"
)
except (ValidationError, ValueError) as exc:
raise RuntimeError("planner_schema_error") from exc
return {"plan": plan.steps, "cursor": 0}
async def execute_node(state: TaskState) -> dict:
step = state["plan"][state["cursor"]]
worker = make_worker(step.tool_budget)
try:
async with asyncio.timeout(step.timeout_seconds):
result = await worker.ainvoke({
"messages": [{
"role": "user",
"content": (
f"总目标:{state['objective']}\n"
f"当前步骤:{step.model_dump_json()}\n"
f"已有证据:{[e.model_dump() for e in state['evidence']]}"
),
}]
})
evidence = result.get("structured_response")
if not isinstance(evidence, Evidence):
raise ValueError("worker 没有返回 Evidence")
except TimeoutError:
evidence = failed_evidence(step, "timeout", "单步执行超过时间预算")
except (ValidationError, ValueError) as exc:
evidence = failed_evidence(step, "schema", str(exc))
except Exception as exc:
# 生产环境还应按网络、鉴权、限流和业务拒绝继续细分。
evidence = failed_evidence(step, "tool_or_model", str(exc))
return {
"evidence": [evidence],
"cursor": state["cursor"] + 1,
}
def after_execute(state: TaskState) -> Literal["execute", "review"]:
return "review" if state["cursor"] >= len(state["plan"]) else "execute"
async def review_node(state: TaskState) -> dict:
review = await reviewer.ainvoke(
f"目标:{state['objective']}\n"
f"计划:{[s.model_dump() for s in state['plan']]}\n"
f"证据:{[e.model_dump() for e in state['evidence']]}\n"
"只按验收条件判断,不要因为文字流畅而放行。"
)
return {"review": review}
def after_review(state: TaskState) -> Literal["replan", "finish"]:
if state["review"] and state["review"].approved:
return "finish"
if state["replan_count"] >= state["max_replans"]:
return "finish"
return "replan"
async def replan_node(state: TaskState) -> dict:
gaps = state["review"].gaps if state["review"] else []
new_plan = await planner.ainvoke(
f"目标:{state['objective']}\n"
f"已有证据:{[e.model_dump() for e in state['evidence']]}\n"
f"未通过原因:{gaps}\n"
"只规划补齐缺口的步骤,不要重复已经通过的工作。"
)
return {
"plan": new_plan.steps,
"cursor": 0,
"replan_count": state["replan_count"] + 1,
}
async def finish_node(state: TaskState) -> dict:
response = await writer.ainvoke(
f"目标:{state['objective']}\n"
f"证据:{[e.model_dump() for e in state['evidence']]}\n"
f"审查:{state['review'].model_dump() if state['review'] else None}\n"
"只根据已有证据成文,并明确仍未解决的缺口。"
)
return {"final_answer": str(response.content)}
builder = StateGraph(TaskState)
builder.add_node("plan", plan_node)
builder.add_node("execute", execute_node)
builder.add_node("review", review_node)
builder.add_node("replan", replan_node)
builder.add_node("finish", finish_node)
builder.add_edge(START, "plan")
builder.add_edge("plan", "execute")
builder.add_conditional_edges(
"execute", after_execute, {"execute": "execute", "review": "review"}
)
builder.add_conditional_edges(
"review", after_review, {"replan": "replan", "finish": "finish"}
)
builder.add_edge("replan", "execute")
builder.add_edge("finish", END)
graph = builder.compile(checkpointer=InMemorySaver())
def choose_max_replans(step_count: int, total_tool_budget: int) -> int:
"""保守示例:预算小就不重规划,任务较长且预算足才允许多一次。"""
if total_tool_budget < step_count * 2:
return 0
return 1 if step_count <= 6 else 2
async def run_task(objective: str, task_id: str) -> TaskState:
initial_state: TaskState = {
"objective": objective,
"plan": [],
"cursor": 0,
"evidence": [],
"review": None,
"replan_count": 0,
"max_replans": choose_max_replans(step_count=8, total_tool_budget=32),
"final_answer": "",
}
return await graph.ainvoke(
initial_state,
config={"configurable": {"thread_id": task_id}},
)
choose_max_replans(...) 只是可解释的保守默认,不是普适公式。真实项目应结合总费用、剩余时长、初始步骤数、失败类型和补缺收益决定。依赖不可用、权限不足这类错误通常不值得重规划;证据漏项才值得。
示例中的失败没有直接让整张图崩溃,而是转成 passed=False 的证据交给 reviewer。这样最终结果可以明确写出“完成了什么、缺了什么、为什么退出”。生产环境还应为 planner 和 reviewer 增加有限重试,并把原始异常写入 trace。
用 Send 真正并行分发
互不依赖的步骤可以并行,例如同时收集 5 家竞品资料。LangGraph 的 Send 会为每个就绪步骤创建一个 worker 输入,Annotated[..., operator.add] 则负责把各分支返回的证据合并起来。
import operator
from typing import Annotated, TypedDict
from langgraph.types import Send
class ParallelState(TypedDict):
objective: str
ready_steps: list[Step]
evidence: Annotated[list[Evidence], operator.add]
class WorkerState(TypedDict):
objective: str
step: Step
def fan_out(state: ParallelState) -> list[Send]:
return [
Send("execute_one", {"objective": state["objective"], "step": step})
for step in state["ready_steps"]
]
async def execute_one(state: WorkerState) -> dict:
step = state["step"]
worker = make_worker(step.tool_budget)
try:
async with asyncio.timeout(step.timeout_seconds):
result = await worker.ainvoke({
"messages": [{
"role": "user",
"content": f"只执行此步骤:{step.model_dump_json()}",
}]
})
evidence = result["structured_response"]
except TimeoutError:
evidence = failed_evidence(step, "timeout", "并行步骤超时")
except Exception as exc:
evidence = failed_evidence(step, "tool_or_model", str(exc))
return {"evidence": [evidence]}
调度器不能把所有步骤一次性扔进 fan_out。它应先计算就绪集合:depends_on 为空,或依赖步骤已经 passed=True 的步骤才能分发。并行分支汇合后,再重新计算下一批就绪步骤。否则只是把有依赖的串行任务错误地同时启动。
异常恢复不能只靠 try/except
生产环境至少要区分四类失败:
- 临时网络错误、限流:有限退避重试;
- 模型结构化输出错误:短重试,仍失败则记录 schema 证据;
- 工具业务拒绝:不要重试,交给 reviewer 或人工;
- 超时、预算耗尽:保存 checkpoint,输出部分结果和缺口。
重试次数也会消耗 tool_budget 和总预算。不能一边说“每步最多 4 次工具调用”,一边在工具外面再悄悄重试 5 次。
什么时候不要使用 Plan-and-Execute
查天气、查单个订单、做一次计算,通常 1 到 3 步就能完成。先生成完整计划反而增加延迟和 token。固定路径则优先用普通 Workflow。
Plan-and-Execute 更适合需要覆盖多个对象、存在依赖、中途可能发现证据缺口,并且最终结果要统一审查的长任务。
小结
Plan-and-Execute 的价值不是让模型先列提纲,而是把长任务变成可调度、可检查、可恢复的状态:
关键结论:计划负责全局覆盖,ReAct 负责局部执行,Harness 约束单次运行,审查器决定是否真的完成。
第六篇会回到架构选型:本篇的 StateGraph,就是“在 Graph 上承载 Plan-and-Execute”的典型案例。
参考资料
- LangGraph, Workflows and agents
- LangGraph, Graph API
- LangGraph, Persistence
- LangChain, Built-in middleware
- LangChain, Structured output
下一篇:如果验收条件已经明确,怎样用外部验证循环让 Coding Agent 持续提交可验证增量?
更多推荐



所有评论(0)