长任务恢复与 Deep Research 续跑
这个教程来自一个真实可运行的 demo:用 LangGraph checkpointer 把一个 7 阶段 Deep Research 流程持久化,演示在 Web UI 里选中某个状态快照后,通过 ResumeRun 沿用同一个 run_id 继续执行,而不是从头重跑。
它不是"假进度条":每个阶段都调用真实外部工具(web_search / web_fetch / LLM / Workspace 写文件),中途可能因为进程重启、网络抖动、用户关闭 Web UI 或主动取消而中断。恢复后必须避免重复调用付费工具、重复写文件、重复写入工作区。
它要解决什么问题
为什么不直接重跑任务
长任务往往包含付费 API、文件写入、工单提交或状态变更。重复执行会产生副作用:重复扣费、重复写文件、重复提交工单。正确做法是:把任务拆成若干"完成后可安全恢复"的业务边界,在每个边界写 LangGraph checkpoint,恢复时只跑 checkpoint 之后的节点,并查询 tool receipt 跳过已完成的工具。
适用场景:
- Agent 执行研究、Deep Research 报告、代码修复、数据分析等耗时任务。
- 任务中途可能因进程重启、网络抖动、用户关闭 Web UI 或主动取消而中断。
- 恢复后必须避免重复调用外部工具(重复扣费、重复写文件、重复写入工作区)。
- 你想先在本地理解 checkpoint / ResumeRun / CancelRun 的工程边界,再接入真实 AgentEngine session store。
图结构
主图由 7 个业务安全点组成,每个阶段完成后由 LangGraph checkpointer 自动写入 StateSnapshot:
| 层级 | LangGraph 用法 | 说明 |
|---|---|---|
| 主图 | StateGraph(ReportState) + checkpoint | 7 个可恢复业务安全点,每完成一个阶段写入 checkpoint。 |
| web_search 子图 | add_conditional_edges(START, fan_out) + Send | 把研究计划里的多个 query 并行分发给 query_agent,汇总候选来源。 |
| 检索补洞循环 | reflection 节点 + _reflection_router | 第一轮证据不足时,按知识缺口追加 gap-fill query 并回到 search_web。 |
| 分析子图 | Send map-reduce | 并行运行多个 reviewer,基于证据表做 LLM 交叉分析。 |
| 条件路由 | add_conditional_edges("analyze_findings", ...) | 默认进入 critic_review;分析阶段模拟一次外部服务失败。 |
| 恢复执行 | graph.astream(None, config=checkpoint_config, stream_mode="updates") | ResumeRun 用 checkpoint config 续跑,不重放 checkpoint 前的节点。 |
7 个业务安全点定义在 stages.py:
REPORT_STAGES = (
ReportStage(key="plan_research", title="规划研究问题", receipt_key="deepresearch:plan:v1", tool_name="llm.plan", ...),
ReportStage(key="search_web", title="检索公开网页", receipt_key="deepresearch:web_search:v1", tool_name="web.search", ...),
ReportStage(key="fetch_sources", title="抓取来源正文", receipt_key="deepresearch:web_fetch:v1", tool_name="web.fetch", ...),
ReportStage(key="screen_evidence", title="筛选证据并去重", receipt_key="deepresearch:evidence_screen:v1", tool_name="evidence.screen", ...),
ReportStage(key="analyze_findings",title="交叉分析发现", receipt_key="deepresearch:analysis:v1", tool_name="llm.analyze", ...),
ReportStage(key="critic_review", title="批判性质检", receipt_key="deepresearch:critic:v1", tool_name="llm.critic", ...),
ReportStage(key="write_report", title="生成研究报告", receipt_key="deepresearch:report:v1", tool_name="workspace.write", ...),
)每个阶段必须有稳定的 receipt_key
receipt_key 是工具幂等键(例如 deepresearch:web_search:v1)。恢复时用它判断某个外部工具是否已经成功执行过。如果 receipt_key 不稳定(例如带时间戳或随机数),恢复时无法匹配,会重复执行工具。建议按 业务对象:工具:版本 生成。
主图编译
主图把 7 个阶段串成顺序链,在 screen_evidence 之后插入 reflection 节点和条件路由:证据充分就进入 analyze_findings,不充分就回到 search_web 补一轮检索。checkpoint 由 LangGraph checkpointer 自动写入,不要自己伪造 checkpoint 本体。
def _compile_graph(checkpointer: Any) -> Any:
graph = StateGraph(ReportState)
for stage in REPORT_STAGES:
graph.add_node(stage.key, _stage_node(stage))
graph.add_node("finalize_report", finalize_report)
graph.add_node("reflection", _run_reflection_node)
graph.add_edge(START, "plan_research")
graph.add_edge("plan_research", "search_web")
graph.add_edge("search_web", "fetch_sources")
graph.add_edge("fetch_sources", "screen_evidence")
graph.add_edge("screen_evidence", "reflection")
graph.add_conditional_edges(
"reflection",
_reflection_router,
{"search_web": "search_web", "analyze_findings": "analyze_findings"},
)
graph.add_conditional_edges(
"analyze_findings",
lambda state: "critic_review" if state.get("requires_critic_review", True) else "write_report",
{"critic_review": "critic_review", "write_report": "write_report"},
)
graph.add_edge("critic_review", "write_report")
graph.add_edge("write_report", "finalize_report")
graph.add_edge("finalize_report", END)
return graph.compile(checkpointer=checkpointer)_reflection_router 控制补洞循环次数,超过 max_research_loops 强制进入分析,避免无限检索:
def _reflection_router(state: ReportState) -> str:
loop_count = int(state.get("research_loop_count") or 0)
cfg = ResearchConfig.from_state(state)
if state.get("is_sufficient") or loop_count >= cfg.max_research_loops:
return "analyze_findings"
return "search_web"子图:用 Send 并行 fan-out
web_search 子图把研究计划里的多个 query 并行分发给 query_agent,每个 query 独立检索,最后由 reduce_queries 汇总。这是 LangGraph 的 map-reduce 模式:add_conditional_edges(START, fan_out) 返回多个 Send,每个 Send 创建一个并行的 query_agent 实例。
def _build_web_search_subgraph():
def fan_out_queries(state: WebSearchState):
queries = state.get("queries") or [state.get("query", "")]
return [Send("query_agent", {"query": state.get("query", ""), "queries": [query]}) for query in queries]
async def query_agent(state: WebSearchState) -> WebSearchState:
query = (state.get("queries") or [state.get("query", "")])[0]
max_results = int(state.get("max_results") or ResearchConfig.from_env().search_results_per_query)
results = await _web_search(query, max_results=max_results)
return {"search_packets": [{"agent": "web_search_agent", "query": query, "results": results}]}
def reduce_queries(state: WebSearchState) -> WebSearchState:
packets = state.get("search_packets") or []
digest = ";".join(f"{packet['query']}={len(packet.get('results') or [])} 条" for packet in packets)
return {"search_digest": digest}
graph = StateGraph(WebSearchState)
graph.add_node("query_agent", query_agent)
graph.add_node("reduce_queries", reduce_queries)
graph.add_conditional_edges(START, fan_out_queries)
graph.add_edge("query_agent", "reduce_queries")
graph.add_edge("reduce_queries", END)
return graph.compile(name="web_search_subgraph")分析子图同理,把证据表并行分发给多个 reviewer(fact_reviewer / implication_reviewer / risk_reviewer),各自独立做 LLM 交叉分析,再汇总:
def fan_out_reviewers(state: AnalysisState):
reviewers = state.get("reviewers") or ["fact_reviewer", "implication_reviewer", "risk_reviewer"]
return [
Send("analysis_reviewer", {
"query": state.get("query", ""),
"evidence_table": state.get("evidence_table", []),
"reviewers": [reviewer],
})
for reviewer in reviewers
]子图可以复杂,主图恢复点要清晰
业务复杂度可以继续加子图(gap-fill 扩展、多 reviewer、多轮反思),但主图恢复点要保持清晰——Web UI 只展示可安全恢复的主图 checkpoint,避免把每个内部节点都暴露成恢复点。子图和工具节点通过 stage_events 转成 tool_call / tool_result 透传给前端。
ResumeRun:从 checkpoint 继续,而不是重跑
这是整个 demo 的核心。ResumeRun 的关键不是"重新跑一遍",而是:
- 根据
run_id和checkpoint_id找到对应的run_checkpoint事件。 - 从
framework_ref.langgraph取出thread_id、可选checkpoint_ns和checkpoint_id。 - 用这些字段构造 LangGraph checkpoint config,以
input=None继续执行。 - 查询 tool receipt,标记已成功的工具调用,避免重复检索、重复写文件、重复写报告。
- 只执行 checkpoint 之后的节点,并继续写入新的 checkpoint。
async for update in graph.astream(
None,
config=config,
context=self.build_native_context(payload.get("platform_context")),
stream_mode="updates",
):
stage = self._stage_from_update(update)
if stage is None:
continue
# 只处理 checkpoint 之后的节点 update,checkpoint 前的节点不会重跑
...input=None 是关键:它表示从 LangGraph checkpoint 继续,而不是重新提交用户问题。如果传了真实的 state,LangGraph 会从 START 重新执行。
构造 checkpoint config 时,thread_id / checkpoint_ns / checkpoint_id 三个字段必须来自 StateSnapshot.config:
@staticmethod
def _checkpoint_config_from(config: dict[str, Any]) -> dict[str, Any]:
configurable = dict((config or {}).get("configurable") or {})
thread_id = str(configurable.get("thread_id") or "").strip()
checkpoint_id = str(configurable.get("checkpoint_id") or "").strip()
if not thread_id:
return dict(config or {})
return {
"configurable": {
"thread_id": thread_id,
"checkpoint_ns": str(configurable.get("checkpoint_ns") or ""),
"checkpoint_id": checkpoint_id,
}
}读取最新 thread 状态时不要带 checkpoint_id
读取"当前 thread 最新状态"必须用不带 checkpoint_id 的 config。如果带上 checkpoint_id,LangGraph 会返回那个历史快照而不是当前续跑状态,导致恢复后读到旧 state。_thread_config_from 专门构造只含 thread_id + checkpoint_ns 的 config 用于读取最新状态。
后台任务:stage 循环与 checkpoint 写入
后台任务在 stream 本体里跑 stage 循环并 yield 事件,由平台消费写 session、写 run_checkpoint 索引、收尾终态。每个阶段执行顺序:tool_call → 执行节点 → 子工具事件透传 → checkpoint → tool_result → 进度文案。
for stage in REPORT_STAGES:
yield {"type": "tool_call", "tool_name": stage.tool_name, ...}
await asyncio.sleep(_stage_delay_seconds())
if _should_fail_stage(stage, "start"):
raise RuntimeError(f"simulated_tool_failure:{stage.key}")
values, framework_ref, checkpoint_id, graph_config = await self._invoke_graph_step(
graph, graph_input=graph_input, config=graph_config, stage=stage, payload=payload,
)
graph_input = None
artifact_path = _artifact_path_for_state(values, stage)
for runtime_event in self._stage_runtime_events(values, stage=stage, run_id=run_id):
yield runtime_event
# checkpoint 先于 tool_result:平台建 tool_receipt 时会从 session 查
# latest checkpoint with matching run_id,必须先写 checkpoint,receipt 才能关联到正确 checkpoint_id。
yield {"type": "checkpoint", "metadata": _checkpoint_metadata_for_ref(stage=stage, run_id=run_id, framework_ref=framework_ref, artifact_path=artifact_path)}
yield {"type": "tool_result", "tool_name": stage.tool_name, "tool_output": {"ok": True, "receipt_key": stage.receipt_key, "artifact_path": artifact_path, "langgraph_checkpoint_id": checkpoint_id}, ...}checkpoint 必须先于 tool_result 写入
平台建 tool_receipt 时会从 session 里查"匹配 run_id 的最新 checkpoint"。如果先写 tool_result 再写 checkpoint,receipt 会关联到上一个阶段的 checkpoint,恢复时工具跳过判断会错位。务必先写 checkpoint 事件,再写 tool_result。
Workspace 写产物
每个阶段都通过 ksadk.toolsets.workspace 写入 Web UI Workspace 的 research/<topic-slug>/ 目录。最终阶段同时写 deepresearch-report.md、deepresearch-report.html、sources.json、evidence.json 和 research-state.json。恢复后通过 tool receipt 和 LangGraph checkpoint 避免重复写文件。
from ksadk.toolsets.workspace import write_workspace_file, write_workspace_files
def _write_workspace_report_files(state, report_markdown, report_html):
base_dir = _workspace_artifact_dir(state)
files = [
{"path": f"{base_dir}/deepresearch-report.md", "content": report_markdown},
{"path": f"{base_dir}/deepresearch-report.html", "content": report_html},
{"path": f"{base_dir}/sources.json", "content": _build_sources_json(state)},
{"path": f"{base_dir}/evidence.json", "content": _build_evidence_json(state, report_markdown)},
{"path": f"{base_dir}/research-state.json", "content": _build_research_state_json(state, report_markdown, report_html)},
]
result = write_workspace_files(files, overwrite=True)
if not isinstance(result, dict) or not result.get("ok"):
raise RuntimeError(f"write_workspace_files failed for {base_dir}: {result}")
written = [str(item.get("path", "")) for item in result.get("written", []) if item.get("path")]
return written, result真实工具边界
这个 demo 默认不是纯 fixture 进度条。Web/API 路径会按下面降级链调用真实工具:
web_search:优先用 ksadk 内置web_search(配KSADK_WEB_SEARCH_PROVIDER=ksyun,复用KSADK_MCP_KEY凭证,走金山云星流 AI 搜索,返回带date的真实时效结果);未配置时回退到 metaso MCP →DEEPRESEARCH_WEB_SEARCH_URL→ 公开搜索页bing,sogou;全部失败才写source=fallback降级结果。web_fetch:首选 ksadkweb_fetch(含 SSRF 防护拦 loopback/内网/metadata + HTML 清洗 + 结果预算),失败时降级到 httpx 直抓。llm.plan/analyze/critic/report:配置OPENAI_API_KEY、OPENAI_BASE_URL、OPENAI_MODEL_NAME后走 OpenAI-compatible/chat/completions。Web/API 真实路径不会用写死文案伪造模型输出。write_workspace_file/write_workspace_files:每个业务安全点都把中间产物写入 Workspace;恢复后通过 tool receipt 避免重复写。
from ksadk.toolsets.workspace import write_workspace_file, write_workspace_files
from ksadk.toolsets.web import web_fetch as _ksadk_web_fetch
from ksadk.toolsets.web import web_search as _ksadk_web_search
async def _web_search(query: str, *, max_results: int = 5) -> list[dict[str, Any]]:
# 优先 ksadk 内置 web_search(ksyun provider),失败再走 metaso → configured_api → 公开搜索降级
ksadk_results = await _search_with_ksadk_web_search(query, max_results=max_results)
if ksadk_results:
return ksadk_results
...为什么直接调用而不绑成 ReAct 工具
本 demo 在 workflow 节点里直接 await _web_search(...),而不是把 web_search 绑成 ReAct 工具让模型自主调用。原因是 Deep Research 的检索流程是确定性的(每个阶段该搜什么由 planner 决定),不需要模型在每轮重新决策"要不要搜"。直接调用保留了多 provider 降级、进程内缓存复用(_KSADK_SEARCH_CACHE)和 receipt 幂等。若你的 agent 需要模型自主选工具,可在图里绑定 ksadk focused 的 web_fetch/web_search 让 ReAct 调用。
CancelRun:协作式取消
CancelRun 建议做成协作式取消,而不是强制 kill:
- API 或 Web UI 只写入取消意图。
- 长任务在每个工具边界检查取消状态。
- 已经开始且不可回滚的工具调用要继续写 receipt,避免下次恢复重复执行。
本 demo 的后台任务把 stage 循环放在 stream 本体里:平台 CancelRun 触发 detached task cancellation,生成器收到 CancelledError 自然退出,平台写 cancelled 终态。stage 失败抛异常,平台写 failed。
不要在原任务仍运行时同时 ResumeRun
真实产品路径应先 CancelRun、等待任务进入 cancelled/failed,再从最新 checkpoint 恢复。否则原任务和恢复任务会并发推进同一个 thread_id,UI 上会出现重复阶段,不能作为验收路径。
运行步骤
- 安装依赖
cd long_task_pg_e2e
uv venv
uv pip install -r requirements.txt
uv pip install -U "ksadk[all]"requirements.txt 锁定了 langgraph、langgraph-checkpoint-postgres、langgraph-checkpoint-sqlite 和 psycopg[binary],本地默认用 InMemorySaver 做交互演示,不访问外部数据库。
- 配置环境变量
cp .env.example .env.env.example 只放占位符,不写真实账号:
OPENAI_API_KEY=your-openai-compatible-api-key
OPENAI_BASE_URL=https://api.openai.com/v1
OPENAI_MODEL_NAME=glm-5.2
# 长任务恢复需要持久化 session/checkpoint。生产环境建议使用 Postgres。
KSADK_SESSION_BACKEND=postgres
KSADK_SESSION_DSN=postgresql://<user>:<password>@<postgres-host>:5432/<database>
KSADK_LANGGRAPH_CHECKPOINT_DSN=postgresql://<user>:<password>@<postgres-host>:5432/<database>
KSADK_SESSION_NAMESPACE=long_task_resume_demo
# 本地演示模式:fixture 用 InMemorySaver;postgres 连真实 PG
KSADK_CHECKPOINT_BACKEND=memory
LONG_TASK_STAGE_DELAY_SECONDS=6真实 token 只放在本地 .env 或运行环境变量中,不提交到仓库。
- 运行行为测试
uv run pytest tests/test_agent_behavior.py -q测试确认普通聊天不伪造 checkpoint、恢复用 LangGraph checkpoint config,而不是重新提交用户问题。
- 启动 Web UI(推荐人工演示)
LONG_TASK_STAGE_DELAY_SECONDS=6 uv run agentengine web .设 6 秒阶段延迟,方便人工在前几个 checkpoint 后观察、继续对话或点击取消。
- 发起研究
调研 2026 年企业采用 AI Agent 平台的关键趋势、主流方案、落地风险和评估建议Agent 会用真实 LLM 流式返回启动确认;后台 run_id、invocation_id、checkpoint_id 只写入 session event 和恢复区,不暴露在普通聊天文本里。
- 观察 checkpoint 与取消恢复
研究工作台会订阅后台 invocation 的会话事件,持续展示 tool_call / tool_result:规划、并行检索、gap-fill query 扩展、网页抓取、证据筛选、多 reviewer 分析、质检和 Workspace 写入。每完成一个业务安全点,才把 LangGraph checkpoint 的 thread_id/checkpoint_id 索引写入 run_checkpoint 事件。
点击"暂停研究"后,前端发送 CancelRun 并等待后台任务停在最近 checkpoint;等待"继续研究"按钮可用后,选择任一 LangGraph 状态快照恢复。ResumeRun 会用该快照里的 framework_ref.langgraph 构造 checkpoint config,以 graph.astream(None, checkpoint_config, stream_mode="updates") 继续执行。
- HTTP E2E 验收(可选)
LONG_TASK_E2E_BASE_URL=http://localhost:9876 uv run python e2e_http.py脚本会流式启动后台 Deep Research → 等待至少两个 checkpoint → CancelRun → GetCheckpointResumePreview → ResumeRun 流式恢复(断言输出含跳过阶段和 deepresearch-report.md)→ ListSessionCheckpoints。
项目结构
agentengine.yaml 声明自定义 runner 类和自研 Web UI:
name: long-task-resume
framework: langgraph
entry_point: agent.py
agent_variable: root_agent
runner_class: agent.LongTaskE2ERunner
ui_profile: custom
ui_path: /research
ui_bundle_path: research-ui/dist
resources:
cpu: "1"
memory: "2Gi"改造成你的业务 Agent
改 REPORT_STAGES 定义自己的业务安全点、_initial_state 构造初始状态、_run_write_report_stage 写自己的报告模板。保留每个阶段的稳定 receipt_key,否则恢复时无法匹配已完成的工具。
REPORT_STAGES = (
ReportStage(key="ingest_data", title="数据接入", receipt_key="myagent:ingest:v1", tool_name="data.ingest", ...),
ReportStage(key="enrich", title="数据富化", receipt_key="myagent:enrich:v1", tool_name="data.enrich", ...),
ReportStage(key="score", title="打分排序", receipt_key="myagent:score:v1", tool_name="ml.score", ...),
ReportStage(key="publish", title="发布结果", receipt_key="myagent:publish:v1", tool_name="api.publish", ...),
)替换 _web_search 或设置 DEEPRESEARCH_WEB_SEARCH_URL(期望 GET ?q=<query>&max_results=<n> 返回 results 列表)。搜索成功后再进入 checkpoint,避免恢复后重复检索。
替换 _web_fetch、_run_fetch_sources_stage。工具输出写入 LangGraph state,再由 checkpointer 持久化,恢复后直接从 state 读取,不重抓。
替换 _call_required_llm、_run_analysis_stage、_run_write_report_stage。Web/API 路径需要真实模型;失败时保留最近 checkpoint,恢复后从当前节点继续,不重放 checkpoint 前的 LLM 调用。
设置 LONG_TASK_RESUME_DEMO_MODE=postgres 和 KSADK_LANGGRAPH_CHECKPOINT_DSN。thread_id / checkpoint_ns / checkpoint_id 必须来自 StateSnapshot.config,不要自己拼。可用 pg_smoke.py 验证 DSN 可连接、可建表、可读写。
保留 request_cancel 和 _run_background_job 的工具边界检查。取消只写意图;不可回滚工具完成后仍要写 receipt,避免下次恢复重复执行。
生产接入建议
生产接入时建议只替换 _web_search、_web_fetch、_run_analysis_stage 和 _run_write_report_stage,保留 REPORT_STAGES 的 receipt key、LangGraph checkpoint config 和 run_checkpoint 事件结构。这样业务逻辑可以换,ResumeRun / ListSessionCheckpoints / CancelRun 的平台闭环不需要重写。
常见问题
checkpoint 应该存在哪里
生产建议放在持久化 session store,并按租户或应用 namespace 隔离。平台写入的 run_checkpoint 只是 UI / 审计索引事件,里面保存 framework_ref.langgraph 和业务摘要;真实 checkpoint 仍由 LangGraph checkpointer 保存。
tool receipt 和日志有什么区别
日志用于观察,receipt 用于恢复决策。恢复逻辑不能只依赖文本日志判断"这个工具跑过没",必须查结构化的 receipt_key + status,否则日志缺失或格式漂移会导致重复执行。
这个 demo 会访问真实平台吗
默认不会。本地默认使用 InMemorySaver 做交互演示,不访问外部数据库;离线 tools.py 的 fixture 用于公开仓库快速展示报告结构。真正可恢复的 LangGraph 路径在 agent.py / runner.py / workflow.py,生产接入请用 agentengine run/web 或部署到目标 pod 后验证。