欢迎光临
我们一直在努力

LangGraph Planning Agent:让 AI 先想清楚再做,中断了还能接着做

目录

一、为什么复杂任务不适合“一把梭”

二、ReAct 与 Planning:不是谁更高级,而是谁更合适

三、核心架构:Plan、Execute、Complete

(一)状态设计:让“步骤”成为可追踪对象

(二)Plan 节点:约束写进 Prompt,解析必须容错

计划质量的三条实用规则

(三)Execute 节点:一次只做一步,前序成果继续递进

1. 前序成果是后续步骤的输入

2. 异常不在节点内部吞掉

(四)断点续跑:为什么输入 None 就能接着做

(五)max_steps:三道防线防止计划失控

(六)图结构:自循环不是细节,而是恢复能力的来源

四、运行演示应该看什么

(一)基本计划执行

(二)步骤数限制

(三)进度追踪

(四)中断恢复

(五)成本模型:为什么断点续跑直接省钱

五、常见问题和线上说明

(一)五个常见坑,以及为什么会踩

1. 坑 1:计划无法终止

2. 坑 2:计划 JSON 解析失败

3. 坑 3:Execute 内部一次跑完整个列表

4. 坑 4:吞掉 LLM 异常

5. 坑 5:原地修改 plan

(二)从 Demo 走向生产:还需要补齐什么

1. 人工审核计划

2. 持久化 Checkpointer

3. 步骤级重试策略

4. 动态重规划

5. 并行执行

6. 可观测性

(三)上线前检查清单

六、FAQ

Q1:Planning Agent 一定比一次性生成效果好吗?

Q2:为什么不把所有步骤一次性并行执行?

Q3:断点续跑会不会重复当前失败步骤?

Q4:failed 状态为什么在 Demo 中没有大量使用?

Q5:Planning 与 Reflection 可以组合吗?

七、总结


干货分享,感谢您的阅读!

这是「LangGraph Agent Engineering Mastery」系列 Stage 4 推理 Agent · 第 2 篇。

这篇文章不只回答“Planning Agent 是什么”,还会把它拆成一个可直接落地的工程模型。读完后,你应该能够:

  • 理解 Planning 与 ReAct 的适用边界,知道什么时候应当“先规划、后执行”;

  • 用 LangGraph 建模 Plan → Execute → Complete 的循环图;

  • 用 current_step_idx、status/result 追踪真实执行进度;

  • 用 max_steps 建立软约束与硬边界,防止计划无限膨胀;

  • 用 Checkpointer 实现中断后的断点续跑,避免重复执行已经完成的 LLM 步骤。

Planning Agent 的价值,不是让模型“多想几步”,而是把一次不可控的大任务,改造成一串可观察、可限制、可恢复、可汇总的小任务。只要每次只推进一步,并在超步结束后持久化状态,流程中断就不再意味着从头再来。

一、为什么复杂任务不适合“一把梭”

想象你接到一个任务:写一份完整的市场分析报告

大多数人不会立刻写第一句话,而是先完成一连串准备动作:确定主题与读者、搜集资料、整理观点、搭建大纲、撰写正文、校验事实、统一表达。这个过程的核心不是“写”,而是先把任务拆明白。

直接把复杂任务一次性交给 LLM,通常会遇到四类问题:

问题具体表现工程后果
容易遗漏关键步骤 模型可能跳过资料搜集,直接产出看似完整的答案 结果结构完整但论证空心
执行顺序混乱 前置分析、验证、汇总没有明确依赖关系 后续步骤缺少可靠输入
中途失败后无法续跑 第 3 步失败,只能从第 1 步重新开始 时间、调用次数与 token 全部浪费
用户看不到进度 黑盒运行时间长,不知道当前卡在哪里 难以审计,也难以干预

Planning 的做法很朴素:把一个不可控的大任务,拆成一串可控的小步骤。 每一步都有描述、状态和结果;每一步结束后都可以持久化;整个流程随时知道“做到哪了”。

二、ReAct 与 Planning:不是谁更高级,而是谁更合适

上一篇的 ReAct 更像“走一步、看一步”:模型观察当前情况,再决定下一次行动。它非常适合搜索、工具调用和探索性任务,因为后续动作往往取决于刚刚获得的新信息。

Planning 则更像一条显式的业务流程:先生成全局计划,再按依赖顺序逐步执行。它更适合报告生成、系统迁移、批量处理、分析工作流等目标明确、步骤相对稳定的任务。

可以用一句话区分:

  • ReAct 解决“下一步该做什么”

  • Planning 解决“整个任务应该怎样有序完成”

两者也可以组合:先由 Planning 给出主计划,再在某个执行步骤内部使用 ReAct 调工具探索。我们示例聚焦于纯 Planning 主流程。

本次教学主要代码前置展示如下:

"""Demo 02: Planning Agent — 任务分解、计划生成、逐步执行、断点续跑。

演示 Planning 推理模式(全链路真实 LLM 调用):
1. Plan 节点:真实 LLM 分析任务,生成结构化执行计划(JSON)
2. Execute 节点:真实 LLM 按计划逐步执行子任务
3. Complete 节点:真实 LLM 汇总各步结果,生成最终交付物
4. 进度追踪:状态中维护 current_step_idx + 每步 status/result
5. 死循环防护:max_steps 限制最大步骤数
6. 断点续跑:Checkpointer 持久化每步进度,真实调用中断(网络异常、
进程崩溃)后,同一 thread_id 传入 None 即可从上次未完成的步骤继续,
已完成的步骤不会重复执行(不浪费已消耗的 LLM 调用)

本 Demo 通过 shared.get_llm() 调用 .env 配置的真实在线模型
(fallback_to_mock=False,需联网)。

运行方式:
python stages/stage4_reasoning/02_planning_agent/main.py
"""

from __future__ import annotations

import json
import re
import sys
from pathlib import Path
from typing import Annotated, TypedDict

from langchain_core.messages import AIMessage, BaseMessage, HumanMessage, SystemMessage
from langgraph.checkpoint.memory import MemorySaver
from langgraph.graph import END, START, StateGraph
from langgraph.graph.message import add_messages

sys.path.insert(0, str(Path(__file__).resolve().parent.parent.parent.parent))

from shared import get_llm, get_logger, log_step, log_success, log_warning

logger = get_logger("demo.04_02_planning_agent")

MAX_PLAN_STEPS = 8

# 故障注入开关(仅供"中断恢复"演示使用):
# 执行到指定步骤 id 时抛出异常,模拟真实 LLM 调用中断(网络故障/进程崩溃)
_FAIL_AT_STEP: int | None = None

# ============================================================
# 真实 LLM(通过 shared.get_llm 获取 .env 配置的在线模型)
# ============================================================
_LLM = None

def _get_planning_llm():
"""获取真实在线 LLM 实例(模块内复用,fallback_to_mock=False 确保真实调用)。"""
global _LLM
if _LLM is None:
_LLM = get_llm(fallback_to_mock=False)
return _LLM

def _parse_llm_json(text: str) -> dict | None:
"""从 LLM 输出中解析 JSON(容忍 Markdown 代码块包裹等格式噪音)。"""
text = text.strip()
if text.startswith("```"):
text = re.sub(r"^```(?:json)?\\s*|\\s*```$", "", text, flags=re.S).strip()
try:
parsed = json.loads(text)
return parsed if isinstance(parsed, dict) else None
except json.JSONDecodeError:
match = re.search(r"\\{.*\\}", text, re.S)
if match:
try:
parsed = json.loads(match.group())
return parsed if isinstance(parsed, dict) else None
except json.JSONDecodeError:
return None
return None

# ============================================================
# Planning State 定义
# ============================================================
class PlanStep(TypedDict):
id: int
description: str
status: str # "pending" | "completed" | "failed"
result: str

class PlanningState(TypedDict):
messages: Annotated[list[BaseMessage], add_messages]
task: str
plan: list[PlanStep]
current_step_idx: int
max_steps: int
is_complete: bool

# ============================================================
# 计划生成(真实 LLM + 解析失败时的保底计划)
# ============================================================
_PLAN_SYSTEM_PROMPT = """你是一个任务规划专家,负责把复杂任务分解为有序的执行计划。

只输出一个 JSON 对象(不要输出任何其他文字):
{"steps": ["步骤1描述", "步骤2描述", …]}

规划要求:
1. 每步描述以动词开头(如"分析/搜索/整理/编写/验证/总结"),具体且可独立执行
2. 步骤之间按依赖关系排序,前面步骤的产出是后面步骤的输入
3. 步骤数量精简,控制在给定上限以内"""

def _fallback_plan_steps(task: str) -> list[str]:
"""LLM 输出无法解析为计划时的保底通用计划(也用于 mock/离线环境)。"""
if "写" in task or "文章" in task or "报告" in task:
return [
"分析写作主题和目标读者",
"搜索相关资料和参考信息",
"整理大纲和关键论点",
"编写正文内容",
"优化文字表达和结构",
"总结核心观点",
]
if "旅行" in task or "规划" in task or "计划" in task:
return [
"分析需求和约束条件",
"搜索可选方案",
"整理方案对比",
"编写详细计划",
"验证计划可行性",
"输出最终方案",
]
if "学习" in task or "教程" in task:
return [
"分析学习目标和前置知识",
"搜索学习资源",
"整理学习路径",
"编写核心概念解释",
"总结学习要点",
]
return [
"分析任务需求",
"搜索相关信息",
"整理关键数据",
"编写解决方案",
"验证方案正确性",
"输出最终结果",
]

def generate_plan_steps(task: str, max_steps: int) -> list[str]:
"""调用真实 LLM 生成计划步骤列表;解析失败时回退到保底计划。"""
llm = _get_planning_llm()
response = llm.invoke([
SystemMessage(content=_PLAN_SYSTEM_PROMPT),
HumanMessage(content=f"任务: {task}\\n步骤数上限: {max_steps}\\n请输出计划 JSON。"),
])
decision = _parse_llm_json(str(response.content))

steps: list[str] = []
if decision:
steps = [str(s).strip() for s in decision.get("steps", []) if str(s).strip()]

if not steps:
log_warning(logger, "LLM 计划输出无法解析为 JSON,使用保底通用计划")
steps = _fallback_plan_steps(task)

return steps

# ============================================================
# 子任务执行器(真实 LLM 调用)
# ============================================================
def execute_subtask(
step_description: str,
task_context: str,
prior_results: list[str] | None = None,
) -> str:
"""调用真实 LLM 执行单个子任务,返回该步骤的执行结果。

注意:这里不捕获 LLM 调用异常——异常向上传播会让本次 invoke 中断,
但已完成步骤的进度已由 Checkpointer 保存,重跑时可从当前步骤继续。
"""
prior_text = (
"\\n".join(f"- 步骤{i + 1}结果: {r}" for i, r in enumerate(prior_results))
if prior_results
else "(这是第一步,暂无前序结果)"
)
prompt = (
f"总任务: {task_context}\\n"
f"前序步骤的成果:\\n{prior_text}\\n\\n"
f"当前要执行的子任务: {step_description}\\n"
f"请基于前序成果完成这个子任务,用中文输出本步骤的具体成果(2~4 句话,"
f"直接给出成果内容,不要复述任务)。"
)
llm = _get_planning_llm()
response = llm.invoke([HumanMessage(content=prompt)])
return str(response.content).strip()

# ============================================================
# Planning 节点实现
# ============================================================
def plan_node(state: PlanningState) -> dict:
"""Plan 节点:真实 LLM 分析任务,生成结构化执行计划。"""
task = state["task"]
max_steps = state["max_steps"]

log_step(logger, "Plan", f"为任务生成执行计划(真实 LLM): '{task[:50]}…'")
print(f"\\n [规划] 分析任务: {task}")

steps = generate_plan_steps(task, max_steps)

if len(steps) > max_steps:
steps = steps[:max_steps]
log_warning(logger, f"计划步骤超出限制,截断为 {max_steps} 步")

plan: list[PlanStep] = [
{"id": i + 1, "description": desc, "status": "pending", "result": ""}
for i, desc in enumerate(steps)
]

print(f" [规划] 生成了 {len(plan)} 步执行计划:")
for step in plan:
print(f" 步骤 {step['id']}: {step['description']}")

return {
"plan": plan,
"current_step_idx": 0,
"is_complete": False,
}

def execute_node(state: PlanningState) -> dict:
"""Execute 节点:真实 LLM 执行当前子任务(每次只推进一步,便于 Checkpoint)。"""
plan = [dict(step) for step in state["plan"]]
idx = state["current_step_idx"]
task = state["task"]

if idx >= len(plan):
log_warning(logger, "所有步骤已执行完毕")
return {"is_complete": True}

current = plan[idx]

# 故障注入(仅"中断恢复"演示):模拟真实 LLM 调用中途失败
if _FAIL_AT_STEP is not None and current["id"] == _FAIL_AT_STEP:
raise RuntimeError(
f"模拟中断:执行步骤 {current['id']} 时 LLM 调用失败(网络异常/进程崩溃)"
)

log_step(logger, "Execute", f"步骤 {current['id']}/{len(plan)}: {current['description']}")
print(f"\\n [执行] 步骤 {current['id']}/{len(plan)}: {current['description']}")

prior_results = [s["result"] for s in plan[:idx] if s["status"] == "completed"]
result = execute_subtask(current["description"], task, prior_results)
current["status"] = "completed"
current["result"] = result

print(f" [结果] {result[:120]}{'…' if len(result) > 120 else ''}")
print(f" [进度] {idx + 1}/{len(plan)} 完成 ({(idx + 1) * 100 // len(plan)}%)")

next_idx = idx + 1
is_complete = next_idx >= len(plan)

return {
"plan": plan,
"current_step_idx": next_idx,
"is_complete": is_complete,
}

def complete_node(state: PlanningState) -> dict:
"""Complete 节点:真实 LLM 汇总所有步骤结果,生成最终交付物。"""
plan = state["plan"]
task = state["task"]

log_success(logger, f"计划执行完毕,共 {len(plan)} 步,调用 LLM 汇总最终结果")
print("\\n [完成] 计划执行完毕,汇总最终结果…")

completed = [s for s in plan if s["status"] == "completed"]
failed = [s for s in plan if s["status"] != "completed"]

results_text = "\\n".join(
f"步骤{s['id']}({s['description']}): {s['result']}" for s in completed
) or "(无已完成步骤)"
prompt = (
f"总任务: {task}\\n"
f"以下是按计划逐步执行后各步骤的成果:\\n{results_text}\\n\\n"
f"请综合以上成果,用中文输出对总任务的最终回答/交付内容(简洁、结构清晰)。"
)
llm = _get_planning_llm()
response = llm.invoke([HumanMessage(content=prompt)])
final_answer = str(response.content).strip()

summary_parts = [f"针对任务「{task}」的最终结果:\\n", final_answer, "\\n— 执行清单 —"]
for step in plan:
status_icon = "✓" if step["status"] == "completed" else "✗"
summary_parts.append(f" {status_icon} 步骤 {step['id']}: {step['description']}")
summary_parts.append(f"\\n总计: {len(completed)} 步完成, {len(failed)} 步未完成")
summary = "\\n".join(summary_parts)

return {
"messages": [AIMessage(content=summary)],
}

# ============================================================
# 条件边
# ============================================================
def should_continue_execution(state: PlanningState) -> str:
"""判断是否继续执行下一步。"""
if state["is_complete"]:
return "complete"
if state["current_step_idx"] >= state["max_steps"]:
log_warning(logger, f"达到最大步骤数 {state['max_steps']},强制结束")
return "complete"
return "execute"

# ============================================================
# 构建 Planning Graph
# ============================================================
def build_planning_graph(max_steps: int = MAX_PLAN_STEPS):
"""构建 Planning Agent 图。

图结构:
START → plan → execute → [execute …] → complete → END

execute 节点每次只执行一步,每个超步(super-step)结束后
Checkpointer 都会保存一次状态——这正是断点续跑的基础。
"""
graph = StateGraph(PlanningState)

graph.add_node("plan", plan_node)
graph.add_node("execute", execute_node)
graph.add_node("complete", complete_node)

graph.add_edge(START, "plan")
graph.add_edge("plan", "execute")
graph.add_conditional_edges(
"execute",
should_continue_execution,
{"execute": "execute", "complete": "complete"},
)
graph.add_edge("complete", END)

return graph

def _initial_state(task: str, max_steps: int) -> dict:
"""构造初始状态。"""
return {
"messages": [HumanMessage(content=task)],
"task": task,
"plan": [],
"current_step_idx": 0,
"max_steps": max_steps,
"is_complete": False,
}

# ============================================================
# 运行演示
# ============================================================
def demo_basic_planning():
"""演示 1:基本计划执行流程。"""
print("\\n— 演示 1: 基本计划执行 —\\n")

graph = build_planning_graph(max_steps=8)
app = graph.compile()

task = "写一篇关于 AI Agent 技术发展的分析报告"
print(f" 任务: {task}")

result = app.invoke(_initial_state(task, max_steps=8))

print(f"\\n {'─' * 50}")
print(" 最终输出:")
print(f" {result['messages'][-1].content}")
return result

def demo_step_limit():
"""演示 2:步骤数限制(死循环防护)。"""
print("\\n— 演示 2: 步骤数限制防护 —\\n")

graph = build_planning_graph(max_steps=3)
app = graph.compile()

task = "规划一次环球旅行,需要很多步骤"
print(f" 任务: {task}")
print(" max_steps = 3 (限制为 3 步)")

result = app.invoke(_initial_state(task, max_steps=3))

completed = [s for s in result["plan"] if s["status"] == "completed"]
print(f"\\n 实际完成步骤: {len(completed)}")
print(f" 步骤限制触发: {'是' if len(completed) <= 3 else '否'}")
return result

def demo_progress_tracking():
"""演示 3:进度追踪。"""
print("\\n— 演示 3: 进度追踪 —\\n")

graph = build_planning_graph(max_steps=6)
app = graph.compile()

task = "学习 LangGraph 框架的入门教程"
print(f" 任务: {task}")

result = app.invoke(_initial_state(task, max_steps=6))

print(f"\\n {'─' * 50}")
print(f" 计划总步骤: {len(result['plan'])}")
print(f" 完成步骤数: {result['current_step_idx']}")
return result

def demo_interrupt_resume():
"""演示 4:中断恢复(Checkpoint 断点续跑)。

模拟真实场景:LLM 调用执行到第 3 步时中断(网络故障/进程崩溃),
由于 Checkpointer 已保存前 2 步的进度,恢复时对同一 thread_id
传入 None 即可从第 3 步继续执行,前 2 步不会重复调用 LLM。
"""
global _FAIL_AT_STEP

print("\\n— 演示 4: 中断恢复(断点续跑) —\\n")

checkpointer = MemorySaver()
graph = build_planning_graph(max_steps=6)
app = graph.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "planning-resume-demo"}}

task = "写一篇关于 LangGraph Checkpoint 机制的分析报告"
print(f" 任务: {task}")
print(" (故障注入: 第 3 步执行时模拟 LLM 调用中断)")

# 阶段 1:执行到第 3 步时人为中断
log_step(logger, "Resume-阶段1", "开始执行,第 3 步将触发模拟中断")
_FAIL_AT_STEP = 3
interrupted = False
result = None
try:
result = app.invoke(_initial_state(task, max_steps=6), config)
except RuntimeError as e:
interrupted = True
log_warning(logger, f"执行中断: {e}")
print(f"\\n [中断] {e}")
finally:
_FAIL_AT_STEP = None # 清除故障(模拟网络恢复/进程重启)

if not interrupted:
# 计划不足 3 步时故障不会触发,直接返回首次执行结果
log_warning(logger, "计划步骤不足 3 步,故障未触发,跳过恢复演示")
return result

# 阶段 2:查看 Checkpoint 中保存的进度
saved_state = app.get_state(config)
saved_plan = saved_state.values.get("plan", [])
completed_before = [s for s in saved_plan if s["status"] == "completed"]
print(f"\\n [Checkpoint] 已保存进度: {len(completed_before)}/{len(saved_plan)} 步完成")
for step in saved_plan:
icon = "✓" if step["status"] == "completed" else "○"
print(f" {icon} 步骤 {step['id']}: {step['description']}")

# 阶段 3:传入 None 从最近一次 Checkpoint 恢复执行
log_step(logger, "Resume-阶段2", "从 Checkpoint 恢复,继续执行未完成步骤")
print("\\n [恢复] 重新运行(输入传 None,同一 thread_id)→ 从第 3 步继续:")
result = app.invoke(None, config)

completed_after = [s for s in result["plan"] if s["status"] == "completed"]
log_success(
logger,
f"断点续跑完成:中断前已完成 {len(completed_before)} 步(未重复执行),"
f"恢复后补齐剩余 {len(completed_after) – len(completed_before)} 步",
)
print(f"\\n {'─' * 50}")
print(f" 中断前完成: {len(completed_before)} 步(结果已持久化,恢复后未重复执行)")
print(f" 恢复后补齐: {len(completed_after) – len(completed_before)} 步")
print(f" 最终状态 : {len(completed_after)}/{len(result['plan'])} 步全部完成")
return result

def run_demo() -> dict:
"""运行 Planning Agent 全部演示。"""
print("=" * 60)
print(" Demo 02: Planning Agent — 任务分解与计划执行(真实 LLM)")
print("=" * 60)

basic_result = demo_basic_planning()
limit_result = demo_step_limit()
progress_result = demo_progress_tracking()
resume_result = demo_interrupt_resume()

print()
print("=" * 60)
print(" 关键概念回顾")
print("=" * 60)
print(" 1. Plan 节点 : 真实 LLM 分析任务,生成结构化执行计划(JSON)")
print(" 2. Execute 节点 : 真实 LLM 逐步执行子任务,前序成果作为上下文")
print(" 3. Complete 节点: 真实 LLM 汇总各步成果,生成最终交付物")
print(" 4. 进度追踪 : current_step_idx + 每步 status/result")
print(" 5. 死循环防护 : max_steps 限制最大步骤数")
print(" 6. 断点续跑 : Checkpointer 保存每步进度,中断后传 None 续跑")
print()

return {
"basic_result": basic_result,
"limit_result": limit_result,
"progress_result": progress_result,
"resume_result": resume_result,
}

if __name__ == "__main__":
run_demo()

三、核心架构:Plan、Execute、Complete

整个 Agent 可以看成三类节点:

  • Plan 节点:LLM 把总任务拆成结构化步骤,并输出 {"steps": […]};

  • Execute 节点:每次只执行当前一步,把前序成果继续注入上下文;

  • Complete 节点:所有步骤结束后,再由 LLM 综合各步成果,生成最终交付物。

  • 这里最重要的设计不是节点数量,而是Execute 每次只执行一步。如果在一个节点里用 for 循环把所有步骤跑完,LangGraph 只会把整个循环视为一个超步。第 3 步中断时,前 2 步也不会形成独立 Checkpoint,恢复能力就失去了意义。

    (一)状态设计:让“步骤”成为可追踪对象

    Planning Agent 的可观测与可恢复,首先来自状态模型。代码中,PlanStep 和 PlanningState 定义如下:

    class PlanStep(TypedDict):
    id: int
    description: str
    status: str # "pending" | "completed" | "failed"
    result: str

    class PlanningState(TypedDict):
    messages: Annotated[list[BaseMessage], add_messages]
    task: str
    plan: list[PlanStep]
    current_step_idx: int
    max_steps: int
    is_complete: bool

    这段结构值得重点关注:

    • plan 不是字符串列表,而是带 id/description/status/result 的对象列表;

    • current_step_idx 明确指出当前执行位置;

    • status 区分 pending、completed、failed;

    • result 保存每一步已经得到的成果;

    • max_steps 把流程上限放进状态,而不是散落在函数内部;

    • is_complete 让条件边可以明确判断是否进入汇总节点。

    你可以把它理解成一个小型工作流状态机。状态不仅告诉系统“下一步做什么”,还保留了“已经做过什么、结果是什么”。这正是断点续跑所需的最小信息集合。

    工程提示:避免原地修改状态 plan 是可变对象。节点里先执行 plan = [dict(step) for step in state["plan"]],再修改副本并返回新对象,可以避免状态历史与 Checkpoint 被意外污染。

    (二)Plan 节点:约束写进 Prompt,解析必须容错

    原始代码先通过系统提示词约束计划格式:

    _PLAN_SYSTEM_PROMPT = """你是一个任务规划专家,负责把复杂任务分解为有序的执行计划。

    只输出一个 JSON 对象(不要输出任何其他文字):
    {"steps": ["步骤1描述", "步骤2描述", …]}

    规划要求:
    1. 每步描述以动词开头(如"分析/搜索/整理/编写/验证/总结"),具体且可独立执行
    2. 步骤之间按依赖关系排序,前面步骤的产出是后面步骤的输入
    3. 步骤数量精简,控制在给定上限以内"""

    这里有三类有效约束:

    • 格式约束:只输出 JSON,不附加解释;

    • 质量约束:步骤以动词开头,具体且能独立执行;

    • 结构约束:按依赖关系排序,并控制在给定上限内。

    但在工程里,不能假设 LLM 永远输出合法 JSON。它可能包一层 Markdown 代码块,也可能在 JSON 前后加解释,甚至返回结构不符合预期。因此,计划生成函数采用“直接解析 → 正则提取 → 保底计划”的容错路径:

    def generate_plan_steps(task: str, max_steps: int) -> list[str]:
    """调用真实 LLM 生成计划步骤列表;解析失败时回退到保底计划。"""
    llm = _get_planning_llm()
    response = llm.invoke([
    SystemMessage(content=_PLAN_SYSTEM_PROMPT),
    HumanMessage(content=f"任务: {task}\\n步骤数上限: {max_steps}\\n请输出计划 JSON。"),
    ])
    decision = _parse_llm_json(str(response.content))

    steps: list[str] = []
    if decision:
    steps = [str(s).strip() for s in decision.get("steps", []) if str(s).strip()]

    if not steps:
    log_warning(logger, "LLM 计划输出无法解析为 JSON,使用保底通用计划")
    steps = _fallback_plan_steps(task)

    return steps

    这段逻辑的价值不是“兼容偶发格式问题”这么简单,而是保证流程的可退化运行:模型输出质量不好时,系统仍能走保底计划,而不是在第一步就完全失败。

    计划质量的三条实用规则

  • 步骤必须有可验证产出:例如“整理三类风险与对应依据”,比“研究风险”更可执行;

  • 步骤粒度要稳定:不要把一个步骤写成一句话,另一个步骤写成完整项目;

  • 依赖关系要可见:后一步应能明确使用前一步结果,而不是多个步骤各自独立写答案。

  • (三)Execute 节点:一次只做一步,前序成果继续递进

    单个子任务执行器会把总任务、前序步骤成果与当前子任务组合成一次真实 LLM 调用:

    def execute_subtask(
    step_description: str,
    task_context: str,
    prior_results: list[str] | None = None,
    ) -> str:
    """调用真实 LLM 执行单个子任务,返回该步骤的执行结果。

    注意:这里不捕获 LLM 调用异常——异常向上传播会让本次 invoke 中断,
    但已完成步骤的进度已由 Checkpointer 保存,重跑时可从当前步骤继续。
    """
    prior_text = (
    "\\n".join(f"- 步骤{i + 1}结果: {r}" for i, r in enumerate(prior_results))
    if prior_results
    else "(这是第一步,暂无前序结果)"
    )
    prompt = (
    f"总任务: {task_context}\\n"
    f"前序步骤的成果:\\n{prior_text}\\n\\n"
    f"当前要执行的子任务: {step_description}\\n"
    f"请基于前序成果完成这个子任务,用中文输出本步骤的具体成果(2~4 句话,"
    f"直接给出成果内容,不要复述任务)。"
    )
    llm = _get_planning_llm()
    response = llm.invoke([HumanMessage(content=prompt)])
    return str(response.content).strip()

    这段代码体现了两个关键原则。

    1. 前序成果是后续步骤的输入

    prior_results 会被整理成“步骤 1 结果、步骤 2 结果……”并注入 prompt。这样,第 4 步不是从空白开始,而是在前 3 步已经得到的事实与结论上继续推进。

    这使线性计划真正具有“递进”关系,而不是把多个独立回答最后简单拼起来。

    2. 异常不在节点内部吞掉

    代码明确说明:LLM 调用异常继续向上传播。本次 invoke 会中断,但此前完成的超步已经由 Checkpointer 保存。

    很多工程师会本能地写 try/except,把异常变成字符串结果并把步骤标记为完成。这样看似“更健壮”,实际上会破坏恢复语义:系统以后会认为该步骤已经成功,最终汇总时却拿到错误信息。

    正确思路是:

    • 节点只负责执行;

    • 异常中断当前调用;

    • 外层决定何时重试;

    • Checkpoint 保留已经成功的进度。

    Execute 节点本身的推进逻辑如下:

    def execute_node(state: PlanningState) -> dict:
    """Execute 节点:真实 LLM 执行当前子任务(每次只推进一步,便于 Checkpoint)。"""
    plan = [dict(step) for step in state["plan"]]
    idx = state["current_step_idx"]
    task = state["task"]

    if idx >= len(plan):
    log_warning(logger, "所有步骤已执行完毕")
    return {"is_complete": True}

    current = plan[idx]

    # 故障注入(仅"中断恢复"演示):模拟真实 LLM 调用中途失败
    if _FAIL_AT_STEP is not None and current["id"] == _FAIL_AT_STEP:
    raise RuntimeError(
    f"模拟中断:执行步骤 {current['id']} 时 LLM 调用失败(网络异常/进程崩溃)"
    )

    log_step(logger, "Execute", f"步骤 {current['id']}/{len(plan)}: {current['description']}")
    print(f"\\n [执行] 步骤 {current['id']}/{len(plan)}: {current['description']}")

    prior_results = [s["result"] for s in plan[:idx] if s["status"] == "completed"]
    result = execute_subtask(current["description"], task, prior_results)
    current["status"] = "completed"
    current["result"] = result

    print(f" [结果] {result[:120]}{'…' if len(result) > 120 else ''}")
    print(f" [进度] {idx + 1}/{len(plan)} 完成 ({(idx + 1) * 100 // len(plan)}%)")

    next_idx = idx + 1
    is_complete = next_idx >= len(plan)

    return {
    "plan": plan,
    "current_step_idx": next_idx,
    "is_complete": is_complete,
    }

    执行完成后,节点会:

    • 将当前步骤标记为 completed;

    • 写入本步 result;

    • 把 current_step_idx 向前推进;

    • 计算是否已经完成全部计划。

    (四)断点续跑:为什么输入 None 就能接着做

    断点续跑成立需要三个条件:

  • 编译图时提供 Checkpointer;

  • 多次运行使用同一个 thread_id;

  • 恢复时传入 None,表示从该线程最近的 Checkpoint 继续,而不是重新初始化状态。

  • 演示代码完整体现了这条链路:

    checkpointer = MemorySaver()
    graph = build_planning_graph(max_steps=6)
    app = graph.compile(checkpointer=checkpointer)
    config = {"configurable": {"thread_id": "planning-resume-demo"}}

    task = "写一篇关于 LangGraph Checkpoint 机制的分析报告"
    print(f" 任务: {task}")
    print(" (故障注入: 第 3 步执行时模拟 LLM 调用中断)")

    # 阶段 1:执行到第 3 步时人为中断
    log_step(logger, "Resume-阶段1", "开始执行,第 3 步将触发模拟中断")
    _FAIL_AT_STEP = 3
    interrupted = False
    result = None
    try:
    result = app.invoke(_initial_state(task, max_steps=6), config)
    except RuntimeError as e:
    interrupted = True
    log_warning(logger, f"执行中断: {e}")
    print(f"\\n [中断] {e}")
    finally:
    _FAIL_AT_STEP = None # 清除故障(模拟网络恢复/进程重启)

    if not interrupted:
    # 计划不足 3 步时故障不会触发,直接返回首次执行结果
    log_warning(logger, "计划步骤不足 3 步,故障未触发,跳过恢复演示")
    return result

    # 阶段 2:查看 Checkpoint 中保存的进度
    saved_state = app.get_state(config)
    saved_plan = saved_state.values.get("plan", [])
    completed_before = [s for s in saved_plan if s["status"] == "completed"]
    print(f"\\n [Checkpoint] 已保存进度: {len(completed_before)}/{len(saved_plan)} 步完成")
    for step in saved_plan:
    icon = "✓" if step["status"] == "completed" else "○"
    print(f" {icon} 步骤 {step['id']}: {step['description']}")

    # 阶段 3:传入 None 从最近一次 Checkpoint 恢复执行
    log_step(logger, "Resume-阶段2", "从 Checkpoint 恢复,继续执行未完成步骤")
    print("\\n [恢复] 重新运行(输入传 None,同一 thread_id)→ 从第 3 步继续:")
    result = app.invoke(None, config)

    恢复时,plan、current_step_idx、每步 status/result 都来自最近一次 Checkpoint。因此:

    • Plan 节点不会重新生成计划;

    • 已完成的步骤不会重复调用 LLM;

    • 当前未完成步骤仍能读取前序结果;

    • 用户看到的进度与中断前保持一致。

    MemorySaver 的边界 本 Demo 使用 MemorySaver,适合本地演示和单进程验证。生产环境需要持久化的 Checkpointer,才能在进程退出或机器重启后继续恢复。核心图结构不需要因此改变,变化主要发生在编译时注入的 Checkpointer 实现。

    (五)max_steps:三道防线防止计划失控

    开放任务很容易被模型拆成几十步,甚至越执行越细。原始实现用了三道防线:

  • Prompt 上限:在 Plan 提示词里明确告诉模型步骤数上限;

  • 计划截断:生成后执行 steps[:max_steps];

  • 条件边硬检查:当 current_step_idx >= max_steps 时强制进入 Complete。

  • 这三层分别解决不同问题:

    • 第一层尽量让模型主动遵守,减少无效输出;

    • 第二层处理模型不遵守约束的情况;

    • 第三层防止状态异常或图循环继续推进。

    工程上,不应只依赖 Prompt。Prompt 的遵循是概率性的,而代码边界是确定性的。

    (六)图结构:自循环不是细节,而是恢复能力的来源

    图的连线非常简洁:

    graph.add_edge(START, "plan")
    graph.add_edge("plan", "execute")
    graph.add_conditional_edges(
    "execute",
    should_continue_execution,
    {"execute": "execute", "complete": "complete"},
    )
    graph.add_edge("complete", END)

    execute → execute 的条件边让每个步骤成为一个独立超步。每次超步完成,Checkpointer 就有机会保存最新状态。

    这意味着中断最多损失“当前这一步”,而不是整个计划。对于真实 LLM 场景,这直接对应成本:已完成步骤的 token 与等待时间不会再次支付。

    四、运行演示应该看什么

    代码包含四组演示,每组都验证一个工程属性。

    (一)基本计划执行

    任务:写一篇关于 AI Agent 技术发展的分析报告。

    观察点:

    • Plan 是否先生成完整步骤;

    • 每一步是否有明确动作与产出;

    • 后续步骤是否使用前序结论;

    • Complete 是否真正综合,而不是机械拼接。

    (二)步骤数限制

    任务:规划一次环球旅行,需要很多步骤;max_steps = 3。

    观察点:

    • 模型是否在 Prompt 约束下主动压缩为 3 步;

    • 即使模型超限,截断是否生效;

    • 条件边是否能保证最终停止。

    (三)进度追踪

    任务:学习 LangGraph 框架的入门教程。

    观察点:

    • current_step_idx 是否稳定增长;

    • 每个步骤状态是否从 pending 变为 completed;

    • 输出中是否能展示 1/6 → 6/6 的进度变化。

    (四)中断恢复

    任务:写一篇关于 LangGraph Checkpoint 机制的分析报告;在第 3 步人为抛出 RuntimeError。

    观察点:

    • 中断后能否读取到前 2 步的状态;

    • 恢复时是否直接从第 3 步开始;

    • 步骤 1、2 是否没有重复调用;

    • 第 3 步是否仍能拿到前两步结果作为上下文。

    (五)成本模型:为什么断点续跑直接省钱

    对于一个 N 步计划,最基础的调用预算是:

    • 1 次 Plan;

    • N 次 Execute;

    • 1 次 Complete。

    因此总调用次数约为 N + 2。例如 8 步计划通常需要 10 次真实调用。

    如果第 7 步失败而系统只能从头开始,前 6 次 Execute 以及 1 次 Plan 都会重复。Checkpointer 的价值因此不只是“体验更好”,而是把恢复粒度从“整个任务”缩小到“当前步骤”。

    五、常见问题和线上说明

    (一)五个常见坑,以及为什么会踩

    1. 坑 1:计划无法终止

    现象:开放任务被拆成大量步骤,执行时间与成本持续增长。 根因:没有显式上限,或只在 Prompt 里写了上限。 处理:Prompt、生成后截断、条件边硬检查三层同时保留。

    2. 坑 2:计划 JSON 解析失败

    现象:json.loads 直接抛错,图在 Plan 节点终止。 根因:LLM 输出带 Markdown 代码块、附加解释或结构错误。 处理:剥离代码块、直接解析、正则提取、保底计划逐层退化。

    3. 坑 3:Execute 内部一次跑完整个列表

    现象:进度从 0% 直接跳到 100%,中断后无法恢复到具体步骤。 根因:所有执行发生在一个超步中。 处理:每次只处理一项,通过条件边自循环。

    4. 坑 4:吞掉 LLM 异常

    现象:失败步骤被错误标记为完成,最终汇总得到垃圾结果。 根因:节点内部捕获异常并继续推进状态。 处理:让异常向上传播,由外层重新调用同一线程恢复。

    5. 坑 5:原地修改 plan

    现象:Checkpoint 历史状态混乱,或 reducer 合并结果难以预测。 根因:直接修改 state 中的可变对象。 处理:复制计划列表与字典,修改副本后返回新状态。

    (二)从 Demo 走向生产:还需要补齐什么

    这个 Demo 已经具备规划、执行、汇总、进度、上限和恢复,但生产系统通常还需要继续增强。

    1. 人工审核计划

    计划生成后先暂停,让用户确认步骤、顺序与风险。高风险动作尤其适合在执行前加入 Human-in-the-Loop。

    2. 持久化 Checkpointer

    内存 Checkpointer 只适合演示。生产中应选择可持久化后端,并保证 thread_id、租户、任务权限和数据保留策略一致。

    3. 步骤级重试策略

    不是所有失败都应无限重试。可以根据异常类型设置:

    • 网络超时:指数退避重试;

    • 内容不合格:进入 Reflection 或 Reviewer 节点;

    • 输入缺失:暂停并请求用户补充;

    • 连续失败:触发重新规划或人工接管。

    4. 动态重规划

    线性计划假设前序结果大体符合预期。如果执行中发现数据缺失或目标变化,需要 Replan 节点只重写剩余步骤,而不是清空已经完成的成果。

    5. 并行执行

    当前示例是线性依赖。如果多个步骤彼此独立,可以建依赖图后并行执行,再汇总到统一节点,减少总等待时间。

    6. 可观测性

    建议记录以下指标:

    指标 用途
    每步耗时 找出最慢步骤与外部服务瓶颈
    每步 token / 成本 评估计划粒度是否过细
    失败次数与异常类型 设计更精确的重试策略
    恢复次数 判断流程稳定性与网络质量
    计划完成率 衡量计划质量与任务可执行性

    (三)上线前检查清单

    • Plan 输出有结构化格式与解析兜底;

    • 计划步骤有明确上限;

    • Execute 每次只推进一个步骤;

    • 每一步都写入 status/result;

    • 状态更新避免原地修改;

    • Checkpointer 已启用并使用稳定的 thread_id;

    • 恢复路径使用 invoke(None, config);

    • 节点内部没有吞掉关键异常;

    • Complete 能区分已完成与未完成步骤;

    • 生产环境已替换为持久化 Checkpointer;

    • 已记录步骤耗时、调用次数与失败原因;

    • 高风险步骤具备人工确认或权限校验。

    六、FAQ

    Q1:Planning Agent 一定比一次性生成效果好吗?

    不一定。简单任务拆步骤反而会增加调用成本。Planning 的优势出现在任务复杂、步骤依赖明确、需要进度与恢复时。

    Q2:为什么不把所有步骤一次性并行执行?

    因为本文示例中的步骤存在前后依赖。只有在确认多个步骤互不依赖时,才适合并行化。

    Q3:断点续跑会不会重复当前失败步骤?

    会重试“当前未完成步骤”,但不会重复已经完成并持久化的步骤。当前步骤如果在 LLM 返回后、Checkpoint 写入前发生进程故障,仍可能再次执行,因此生产系统还应考虑幂等性。

    Q4:failed 状态为什么在 Demo 中没有大量使用?

    当前策略选择“异常中断 + 恢复重试”,因此失败不会被吞掉并写成最终状态。需要支持跳过、人工处理或继续汇总时,可以显式写入 failed 并扩展条件边。

    Q5:Planning 与 Reflection 可以组合吗?

    可以。常见做法是每个步骤完成后先进入 Reviewer,质量不达标就修订;全部步骤通过后再进入 Complete。这也是下一阶段可继续扩展的方向。

    七、总结

    Planning Agent 的本质,是把“大任务的一次性生成”变成“有状态的小步骤执行”。真正让它具备工程价值的,不只是 LLM 会列计划,而是以下四件事同时成立:

  • 计划显式存在于状态中;

  • Execute 每次只推进一个步骤;

  • max_steps 提供确定性的终止边界;

  • Checkpointer 在超步之间保存进度。

  • 做到这些之后,Agent 才从一个长时间运行的黑盒,变成一个可观察、可限制、可恢复的工作流。

    下一篇可以继续进入 Reflection Agent:让模型不只“按计划做完”,还会对结果进行评审、修订和迭代,直到达到质量阈值。

     系列导航:LangGraph从零构建生产级 AI Agent 平台的递进式学习项目-CSDN博客

    赞(0)
    未经允许不得转载:171主机测评 » LangGraph Planning Agent:让 AI 先想清楚再做,中断了还能接着做
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址