KsADK
教程最佳实践

长任务恢复与 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) + checkpoint7 个可恢复业务安全点,每完成一个阶段写入 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:

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 本体

workflow.py
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 强制进入分析,避免无限检索:

workflow.py
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 实例。

workflow.py
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 交叉分析,再汇总:

workflow.py
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 的关键不是"重新跑一遍",而是:

  1. 根据 run_idcheckpoint_id 找到对应的 run_checkpoint 事件。
  2. framework_ref.langgraph 取出 thread_id、可选 checkpoint_nscheckpoint_id
  3. 用这些字段构造 LangGraph checkpoint config,以 input=None 继续执行。
  4. 查询 tool receipt,标记已成功的工具调用,避免重复检索、重复写文件、重复写报告。
  5. 只执行 checkpoint 之后的节点,并继续写入新的 checkpoint。
runner.py
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:

runner.py
@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 → 执行节点 → 子工具事件透传 → checkpointtool_result → 进度文案。

runner.py
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.mddeepresearch-report.htmlsources.jsonevidence.jsonresearch-state.json。恢复后通过 tool receipt 和 LangGraph checkpoint 避免重复写文件。

tools.py
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:首选 ksadk web_fetch(含 SSRF 防护拦 loopback/内网/metadata + HTML 清洗 + 结果预算),失败时降级到 httpx 直抓。
  • llm.plan/analyze/critic/report:配置 OPENAI_API_KEYOPENAI_BASE_URLOPENAI_MODEL_NAME 后走 OpenAI-compatible /chat/completions。Web/API 真实路径不会用写死文案伪造模型输出。
  • write_workspace_file / write_workspace_files:每个业务安全点都把中间产物写入 Workspace;恢复后通过 tool receipt 避免重复写。
tools.py
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 上会出现重复阶段,不能作为验收路径。

运行步骤

  1. 安装依赖
cd long_task_pg_e2e
uv venv
uv pip install -r requirements.txt
uv pip install -U "ksadk[all]"

requirements.txt 锁定了 langgraphlanggraph-checkpoint-postgreslanggraph-checkpoint-sqlitepsycopg[binary],本地默认用 InMemorySaver 做交互演示,不访问外部数据库。

  1. 配置环境变量
cp .env.example .env

.env.example 只放占位符,不写真实账号:

.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 或运行环境变量中,不提交到仓库。

  1. 运行行为测试
uv run pytest tests/test_agent_behavior.py -q

测试确认普通聊天不伪造 checkpoint、恢复用 LangGraph checkpoint config,而不是重新提交用户问题。

  1. 启动 Web UI(推荐人工演示)
LONG_TASK_STAGE_DELAY_SECONDS=6 uv run agentengine web .

设 6 秒阶段延迟,方便人工在前几个 checkpoint 后观察、继续对话或点击取消。

  1. 发起研究
调研 2026 年企业采用 AI Agent 平台的关键趋势、主流方案、落地风险和评估建议

Agent 会用真实 LLM 流式返回启动确认;后台 run_idinvocation_idcheckpoint_id 只写入 session event 和恢复区,不暴露在普通聊天文本里。

  1. 观察 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") 继续执行。

  1. 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
requirements.txt
.env.example
agent.py
runner.py
workflow.py
tools.py
stages.py
config.py
prompts.py
llm_client.py
e2e_http.py
pg_smoke.py

agentengine.yaml 声明自定义 runner 类和自研 Web UI:

agentengine.yaml
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,否则恢复时无法匹配已完成的工具。

stages.py
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=postgresKSADK_LANGGRAPH_CHECKPOINT_DSNthread_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 后验证。

后续阅读

本页导航