构建可恢复的 AI 图流水线
我正在构建 Intelligent Book Summarizer(IBS),一个 AI 驱动的应用,旨在把整本书和长篇文档变成结构化、有用的知识。
IBS 不是从一大块文本生成一份笼统的摘要,而是使用多阶段处理流水线。它首先提取并规范化文档,识别其结构,把它拆成可管理的块,逐步总结它们,然后把那些结果聚合为更高级别的输出。
该系统可以从 PDF、EPUB、Markdown 和纯文本等格式生成高管摘要、章节级洞察、关键概念、学习笔记和术语表。
从技术角度看,IBS 结合了基于 React 的前端、Python/FastAPI 后端和一个 LangGraph 编排的 LLM 工作流。
该流水线围绕模块化、幂等的处理步骤设计,允许长任务从中间产物恢复,而不是从头重启。它还兼容 OpenAI 兼容的模型端点,使 AI 层相对灵活,且与单一提供商解耦。
这个项目有趣的地方不在于简单地"摘要一本书",而在于设计一个可靠的工作流来处理长文档,同时跨多个阶段保留结构、上下文和逐步精炼的信息。
1、长流水线会失败,而重跑并不免费
一本书长度的输入会被切成大约两百块,每一块在合并之前都有自己的模型调用。
我第一次在真实文档上端到端运行时,它死在了第一百三十多块之后的某个地方。
那天失败的原因是速率限制。其他日子可能是请求超时、部署期间的容器重启、进程因内存被杀,或者一个格式错误的块在进入提示词时抛错。这份清单很无聊,而这正是重点:这些都不是什么稀奇事。任何运行二十分钟的流水线都会撞上其中好几个。
显而易见的回应是再跑一次。两百次调用的情况下,这个回应错得离谱,而且是三重错误。
Token 要付两次钱。 墙钟时间要花两次。而且这次运行会径直走回杀死第一次尝试的那个速率限制,因为请求模式一点都没变。重试并不能恢复这次运行,它只是重复它。
存在一个阈值,低于它这一切都无所谓。如果整个运行只需要二三十秒,可恢复性纯粹是开销。上面讲的一切都适用于那条线之上——在那里,一次失败要花真金白银和真正的分钟数,而"从头开始"和"继续下去"之间的区别,就是一个可用的服务和一次演示之间的区别。
2、这个模式来自的应用
Intelligent Book Summarizer 接收 PDF、EPUB、TXT 或 Markdown 文件,产出一份最终摘要、逐章摘要、关键概念、学习笔记和一份术语表,可下载为 Markdown、HTML、PDF 或 JSON。
选择一种摘要风格——高管摘要、深度技术、学习、逐章——再选一种目标语言,上传一本书,流水线就完成剩下的工作。
这个技术栈是有意保持普通的。一个 React + Vite 前端通过 REST 和轮询与 FastAPI 后端通信;后端拥有全部 AI 逻辑,并把摘要作为** LangGraph** 工作流运行在任何 OpenAI 兼容的端点上。输出落在文件系统上而不是数据库里,这一点对后面讲的内容至关重要。
那张图里的两个决定对本文的其余部分至关重要。
模型被藏在薄薄的一层隔离层后面,所以当没有配置 API 密钥时,整个应用运行在一个确定性的 mock 上——这就是测试如何保持免费且封闭。而输出是每次运行专属目录里的文件,不是数据库行。第二点正是让下面的模式成为可能的原因。
这个应用是一个个人项目,它的仓库是私有的。与其把你送到你读不到的代码面前,我把这个模式重建为一个随本文附带的小型可运行示例——同样的想法,只是剥离了产品外壳。
3、模式:产物优先、幂等的节点
基本规则很简单:每个节点先检查自己的输出产物是否已经存在于磁盘上,如果存在就提前返回。
这是一种与大多数工作流框架所鼓励的不同的心智模型。通常的做法是检查点框架的状态:每一步之后,序列化引擎所知道的东西,以便之后的一次运行可以重新水合。产物优先的方法把这个倒过来了。
输出就是检查点。提取出的文本、分块文件和摘要本来就要写进磁盘,因为它们是提供给用户的内容,而它们的存在就是这一步完成的证明。
由此带来的结果说服了我。没有"恢复代码路径"。启动一次运行的函数和恢复一次运行的函数是同一个函数,以同样的方式调用,使用同样的参数。
在摘要器里,POST /jobs/{id}/start 和 POST /jobs/{id}/resume 调用的是完全相同的方法;resume 端点唯一的不同是往事件日志里写一行"resume requested"。
这两个想法的先后顺序很重要。幂等性是对每个节点施加的要求。
可恢复性是从中涌现出来的属性。你不是构建可恢复性;你构建幂等的节点,然后发现你已经拥有了它。
这就是摘要器里真正的工作流:按固定顺序排列的十一个节点,从文本提取、分块、逐块摘要,一直到最终导出。
副标题就是整个设计——每个节点先检查自己的输出产物。把图里的文件名当作示意即可;下面所有具体的路径都来自那个可运行示例。
4、一个最小化的可恢复图
示例位于:https://github.com/gandrenacci/resumable-langgraph
这是一个基于纯文本文件的五节点 LangGraph 流水线。
加载文档、分块、总结每个块、可选地提取关键点、导出报告。它无需 API 密钥、零成本即可离线运行,因为"模型"是一个确定性的桩,它会数自己的调用次数,并且可以被要求按需崩溃。
图本身是项目里最无趣的文件,而且这是刻意为之的。五个节点从一个有序列表连成一条直线。
没有条件边、没有检查状态并向前跳转的路由器、没有恢复分支。
图永远不知道一次运行是一次恢复——每一个跳过决策都住在节点内部,所以它不会与拓扑结构脱节。真正的摘要器有十一个节点而不是五个,除此之外是同一个文件。
5、产物优先的节点长什么样
形状是:守卫、提前返回、干活、写入、标记完成。
# code/resumable/nodes.py
def load_document(self, state: dict[str, Any]) -> dict[str, Any]:
"""Copy the source document into the run directory."""
run_id = state["run_id"]
paths = self.store.paths(run_id)
if paths.raw_text.exists() and paths.raw_text.stat().st_size > 0:
self._skip(run_id, "load_document")
return {"raw_text_path": str(paths.raw_text)}
self._begin(run_id, "load_document", "loading source document")
text = Path(state["input_path"]).read_text(encoding="utf-8")
if not text.strip():
raise ValueError("source document is empty")
paths.raw_text.write_text(text, encoding="utf-8")
self._register(run_id, "raw_text", paths.raw_text)
self._done(run_id, "load_document")
return {"raw_text_path": str(paths.raw_text)}
有两行承载着这个模式,而且它们都容易在微妙之处出错。
守卫是 exists() and stat().st_size > 0,而不是 exists()。
写入中途崩溃会留下一个零字节文件,而一个单纯的存在性检查会把它当作已完成的工作:运行继续越过一个什么也没产出的步骤,然后在后面的某个让人困惑的地方失败。
提前返回必须产生与工作路径相同的状态增量。两个分支都返回 {"raw_text_path": ...}。
分块节点是这一点代价高昂的地方,因为在跳过路径上,它要重新读取自己的清单,纯粹是为了重建 chunks 列表。跳过这一步,恢复后的运行在到达摘要步骤时状态里没有块,就会死在一个看起来和实际错误完全不像的 KeyError 上。
并非每个节点都配拥有一个产物。 摘要器的结构检测节点根本没有守卫:它是纯内存的,几微秒就完成。缓存它,然后永远维护那个缓存的模式,成本比重算它更高。
要点:
- 守卫要检查真正的完整性(
exists()且非空),而不是路径存在。 - 跳过路径必须返回与工作路径完全相同的状态增量。
- 只缓存昂贵的东西。廉价的纯节点应该直接重跑。
6、一鱼两吃的状态文件
产物告诉你什么完成了,而不是现在正在发生什么,而一个盯着进度条的用户需要后者。
所以每次运行在 state/job_state.json 有一个 JSON 文件,保存状态、当前步骤、已完成的步骤、分块计数器、选项标志、一个产物注册表和一个有界事件日志。
进度在分块摘要时给部分学分,否则一个运行八分钟的节点会显示一条冻结的进度条,用户会以为运行死了。
产物注册表是整个想法最具体的微小表达。它只列出磁盘上实际存在的文件,而这份清单驱动着 UI 中哪些下载按钮被启用。没有任何数据库列记录关键点是否已生成。文件就是真相。
每次恢复时都会先读这个文件,所以一次撕裂的写入是致命的。
# code/resumable/store.py
def save_state(self, state: RunState) -> None:
state.updated_at = _now()
paths = self.paths(state.run_id)
paths.state_dir.mkdir(parents=True, exist_ok=True)
tmp = paths.state_file.with_suffix(".json.tmp")
tmp.write_text(json.dumps(state.to_dict(), indent=2), encoding="utf-8")
os.replace(tmp, paths.state_file)
os.replace 是承重的那次调用:在 POSIX 上,对已有路径的重命名是原子的,所以观察者看到的是旧文件或新文件,永远不会看到混合物。
原地写入会先截断,而如果那时崩溃,就会留下下一次恢复无法解析的 JSON。每次运行一把锁解决了第二种交错来源,因为 API 层和节点对同一个文件执行读-改-写循环。
要点:一个类拥有所有状态变更,通过临时文件写入,并在写入期间持有锁。它对单进程、单一本地文件系统是正确的,而且它不是分布式锁。
7、困难的案例:在节点内部恢复
节点级恢复是入场券,但它解决不了我实际遇到的问题。摘要节点为每个块做一次模型调用——单个节点内部两百次调用。
如果恢复的粒度是节点,那么在第 137 次调用时崩溃会丢弃 136 次成功的调用并从零开始。这条流水线名义上是可恢复的,却什么都没恢复出来。
所以这个循环需要自己的持久化记录,记录哪些条目已经完成。
# code/resumable/nodes.py (trimmed: event/progress calls elided)
existing: list[dict[str, Any]] = []
if paths.chunk_summaries.exists():
existing = json.loads(paths.chunk_summaries.read_text(encoding="utf-8"))
done_indices = {item["index"] for item in existing}
summaries = list(existing)
# Seed continuity from the last completed summary, not from an empty
# string, or the narrative breaks exactly at the crash boundary.
previous = summaries[-1]["summary"] if summaries else ""
for chunk in chunks:
if chunk["index"] in done_indices:
continue
text = (paths.chunks_dir / chunk["file"]).read_text(encoding="utf-8")
summary = self.llm.complete(text, previous)
previous = summary
summaries.append({"index": chunk["index"], "summary": summary})
summaries.sort(key=lambda item: item["index"])
# Persist after every chunk. N small writes cost microseconds each;
# the call they protect costs seconds and money.
paths.chunk_summaries.write_text(
json.dumps(summaries, indent=2), encoding="utf-8"
)
这个部分结果文件兼任这个循环的恢复清单。跳过键是块索引,它跨运行保持稳定,因为它来自分块节点写入的清单。任何从易变的东西(比如重新生成的列表里的位置)推导出的键,都会悄悄重做或跳过错误的条目。
在循环内写入是人们会退缩的交易:N 次写入而不是一次。它根本不是问题。为了省 I/O 而批量写入,意味着丢失自上次 flush 以来的每一条摘要——而这恰恰是这个模式存在要保护的工作。
连续性种子是我第一个搞错的点。每条摘要都以上一条为条件,这样输出读起来是连续的散文,而不是两百段互不关联的段落。
我的第一个版本初始化 previous = "" 并在循环运行时填充它:全新运行时正确,每次恢复时错误。第 137 块会被总结得好像它前面什么都没有,而文本正好在崩溃发生的地方猛地一滞。没有完整性测试能抓住这个,因为输出是完整的。它只是更差了。
Token 和成本计数器有同样的形状。如果一个节点写它运行的总量,一次恢复的运行会重新记录已经算过的开销,而仪表盘每次有什么失败都会膨胀。写增量。
**要点:**节点内恢复需要三样东西——一个边走边写的逐条目产物、一个跨运行稳定的跳过键、以及记录为增量而非总量的累计计数器。
8、看着它崩溃并恢复
这个演示在第 4 次模型调用时崩溃一次运行,然后调用同一个入口点:没有恢复标志、没有起始偏移、没有单独的函数。两次运行之间,它把第一个块的摘要覆盖为一个模型永远不可能产生的哨兵字符串。
如果那个字符串存活到了最终输出里,就无可辩驳地证明那个块没有被重做——比数调用次数更有说服力,因为它落在用户会下载的产物里。
# uv run python demo_resume.py (trimmed: repeated chunk lines elided)
=== RUN 1: crashes on the 4th LLM call ===
[load_document] loading source document
[chunk_document] created 8 chunks
[summarize_chunks] chunk 1/8 summarized
[summarize_chunks] chunk 2/8 summarized
[summarize_chunks] chunk 3/8 summarized
[summarize_chunks] failed: provider failed on call 4
status : failed
chunks done : 3/8
progress : 0.475
What survived on disk:
chunks/manifest.json (667 bytes)
input/raw.txt (9954 bytes)
state/job_state.json (1533 bytes)
summaries/chunk_summaries.json (377 bytes)
Planted sentinel in chunk 0: SENTINEL-DO-NOT-REGENERATE
=== RUN 2: same entrypoint, no arguments changed ===
[load_document] skipped - artifact already present
[chunk_document] skipped - artifact already present
[summarize_chunks] resuming: 3/8 chunks already summarized
[summarize_chunks] chunk 4/8 summarized
...
[export_report] run completed
status : completed
LLM calls run 2 : 5 (of 8 chunks)
Sentinel survived: the first chunk was never re-summarized.
八个块,三个已完成,五次调用。progress: 0.475 是五个步骤中的两个完成,加上第三个完成了八分之三:(2 + 3/8) / 5。
失败后的目录清单是我觉得最让人安心的部分。运行产生的一切都还在那里,原封未动——失败处理器记录错误、设置状态,然后什么都不删。失败后清理的诱惑很强,但那些残骸就是检查点。
9、第二个回报:给一次完成的运行添加输出
几个月后来了一个与失败无关的新需求。我提交了一份关闭了关键点的文档,等运行完成,读完摘要,然后决定我终究还是想要关键点。
当昂贵的那部分已经躺在磁盘上时,为了多一个产物就把一切重跑一遍是荒谬的。
我把它估计为真材实料的工作:一个部分再生成入口点、一种表达要执行图的哪个子集的方式、它自己的状态跟踪、它自己的一堆 bug。我刚开始画草图,就注意到流水线已经做了这件事。
实际需要的:在状态文件里设置标志,删除那个必须重建的组合输出,通过同一个运行器重新提交同一次运行。每个已有产物的节点都会跳过,新启用的节点运行,导出节点重建。第一次八次模型调用,第二次一次——而且新的下载出现了,没有任何人更新标志来声明这件事,因为注册表是从磁盘上的东西重建的。
做新工作的节点有一个双重守卫:先检查选项标志,然后是普通的产物检查。它们彼此不知道对方的存在却能组合,这就是为什么在完成的上翻转标志是可行的。
选项守卫现在通过了,产物守卫什么也没找到,于是节点运行。
有一个设计细节让这一切成立,而它正是评审时看起来不对的那个。导出节点没有产物守卫。它总是重跑,并从磁盘上已有的东西重建。
给它加守卫本会更一致,但那会完全毁掉这个功能,因为报告在第一次运行时就已经存在,节点会跳过,而新的关键点永远到不了用户手里。归约(reduce)步骤应该重跑,因为它的输入在它下面变化。
要点:当你为一个原因构建的属性把一项计划中的功能变成了一个配置变更时,这就是抽象切对了位置的可靠证据。
10、这要付出什么代价
每个中间产物都是一个文件,而文件会累积。需要一个保留策略,而没有人会在第一天写一个,所以麻烦的第一个信号通常是一条磁盘告警。
更重的代价是路径和模式变成了接口:一旦守卫检查 summaries/chunk_summaries.json,那个路径就是 API。
在不变的路径背后改变模式是危险的情况,因为守卫说完成了,而下游节点读到一个它不理解的结构的文件。
缓存失效是未解决的部分。 改变提示词、模型、块大小或输入,磁盘上的每一个产物都已过期,却仍然看起来完整。守卫会把它们全部跳过,然后递给用户一份由已经不存在的提示词生成的摘要构建出的报告。
这在我们的实现里没有解决,而我宁愿说清楚,也不愿把一个设计描述得好像已经上线了。我会构建的修复是一个指纹:对影响节点输出的配置做哈希,把不匹配当作"未完成"。别扭的地方在于失效必须是逐节点的。
一次恢复的运行还会混合两次模型调用(可能相隔数天)的输出。对摘要来说这没问题。但当整个输出内部的连贯性就是产品本身时,就不行了。
LangGraph 自己的检查点器很好,如果你需要对话记忆、时间旅行或人在环路的打断,用它们,而不是重新发明一个更差的版本。
我对它们的反对是刻意狭窄的:检查点器在图的步骤之间持久化状态,所以它本身并不能把一次运行从某个节点内部的两次调用循环中救出来。把这个循环重构为扇出(fan-out)可以做到,但要承担这个模式一直在回避的编排复杂度。而且如果一次运行必须挺过它启动所在那台机器的丢失,你需要一个持久化的共享存储层,而不是进程本地的文件系统状态。
当检查点本来就是必须被生产并交付的产物时,产物优先的可恢复性就赢了:唯一的新代码是每个节点顶部的一个守卫。当中间产物本身不值得保留时,这个模式就是一个带保留策略和失效问题的缓存,而计算看起来就非常不同了。
11、两个可以带走的模式
这篇文章可以归结为两个可复用的模式。
11.1 架构模式
对于长运行的 LLM 工作流:
长运行 LLM 流水线 → 持久化产物 → 幂等执行 → 可恢复性 → 节点内恢复
关键想法是把中间输出当作持久的执行状态。当有效产物已经存在时,昂贵的步骤可以跳过工作,而长运行的循环可以增量地持久化部分进度。于是可恢复性成为流水线设计方式的一个结果,而不是一条单独的恢复路径。
11.2 实现模式
一种实践上验证这个架构的方式是:
真实失败 → 可恢复性模式 → 最小实现 → 节点内失败 → 崩溃并恢复的演示 → 扩展一次完成的运行 → 明确的限制
这个序列在这个具体示例之外也很有用。从一次重要的失败开始,隔离出最小的可复用机制,在真实中断下测试它,验证已完成的工作可以被复用,最后把权衡和限制说清楚。
我出发时是为了不再为一次失败的运行付两次钱,而最终我得到了一条"再做一次工作"和"做工作"是同一个指令的流水线。运行示例的代码在 code/ 下,运行它不需要任何凭据。
原文链接: Building an AI Resumable Graph Pipeline
汇智网翻译整理,转载请标明出处