文章摘要: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”的典型案例。

参考资料

  1. LangGraph, Workflows and agents
  2. LangGraph, Graph API
  3. LangGraph, Persistence
  4. LangChain, Built-in middleware
  5. LangChain, Structured output

下一篇:如果验收条件已经明确,怎样用外部验证循环让 Coding Agent 持续提交可验证增量?

Logo

小龙虾开发者社区是 CSDN 旗下专注 OpenClaw 生态的官方阵地,聚焦技能开发、插件实践与部署教程,为开发者提供可直接落地的方案、工具与交流平台,助力高效构建与落地 AI 应用

更多推荐