BaseCheckpointSaver:持久化介面的形狀
職責
在 LangGraph 的執行時裡,Pregel 引擎每跑完一個超步 (superstep) 都要把自己當前所有通道 (channel) 的值連同這一步產生的寫 (writes) 一起存下來,這樣下次進來同一個 thread_id 才能接著上一次的狀態往下跑,也能在被打斷 (interrupt) 之後恢復、回放歷史、做時間旅行除錯。BaseCheckpointSaver 就是這層「持久化」的抽象介面 (BaseCheckpointSaver:176)。它本身不關心你把資料存到哪——記憶體 dict、SQLite、Postgres、Redis 都行——只規定 saver 必須實作的幾個方法:put / get / get_tuple / list / put_writes (put / get / list / put_writes:227-318)。
它同時把檢查點 (checkpoint) 這個概念的形狀釘死成 Checkpoint TypedDict:Checkpoint TypedDict:92-123,裡面裝著 channel_values、channel_versions、versions_seen、id / ts / v 這幾樣。put 收到的不只是一個 Checkpoint,還有 metadata 和 new_versions,saver 決定哪些欄位進主表、哪些進 blob 表。put_writes 則單獨存「這條任務在這步產出的寫」,跟 checkpoint 解耦,這樣寫可以增量落盤,不用每寫一次就重寫整份快照。
換句話說,BaseCheckpointSaver 是 PregelLoop 和具體儲存後端之間的合約:引擎只認這套介面,後端換實作不用改引擎;後端寫自己的 saver 只要照著五個方法填,引擎就自動獲得該後端的並行和持久化特性。
設計動機
- 狀態和寫分開存:
put存整份快照,put_writes存單個 task 在某步內的增量寫 (put_writes:300-318)。拆開是為了讓中斷恢復時能精確還原「這一步哪個 task 已經寫過、哪個還沒」,而不是只能整體回滾到上一個 checkpoint。 thread_id當主鍵:put回傳的 config 裡只裝thread_id/checkpoint_ns/checkpoint_id(put:277-298),跨呼叫、跨行程都以這三個欄位定位一個 checkpoint。checkpoint_ns是給子圖用的命名空間,空串就是頂層圖。- 同步 + 非同步雙套介面:每個寫方法都有
aput/aget/alist/aput_writes的 async 版本 (async methods:417-509),預設實作拋NotImplementedError。saver 子類別可以只實作同步版,非同步版預設走asyncio.to_thread轉發,也可以自己寫真正的 async 路徑走原生驅動。 - 元資料帶來源標籤:
CheckpointMetadata裡source欄位取input/loop/update/fork(source:41-48)。list 和過濾都靠 metadata,所以 saver 必須把 metadata 跟 checkpoint 一起存,不能丟。 - DeltaChannel 旁路:
prune/delete_for_runs/copy_thread的 docstring 反覆警告 (DeltaChannel-aware ops:320-415),如果圖用了DeltaChannel,自訂 saver 不能只刪最新 checkpoint,得把祖先鏈上的checkpoint_writes一併保留,否則 delta 重建會靜默回傳空。
關鍵檔案
BaseCheckpointSaver 類別:176— 抽象基類本體,規定 put/get/list/put_writes 等介面。Checkpoint TypedDict:92-123— 一個 checkpoint 的欄位形狀:channel_values/channel_versions/versions_seen/id/ts。CheckpointTuple:139-146—get_tuple回傳值,把 checkpoint、metadata、parent_config、pending_writes 打包成一個 NamedTuple。CheckpointMetadata:38-86—source/step/parents/run_id/ DeltaChannel 用的counters_since_delta_snapshot。get 方法:227-237—get預設委託get_tuple,只回傳checkpoint欄位。put 方法:277-298— 存整份 checkpoint + metadata,回傳更新後的 config。put_writes 方法:300-318— 存 task 在某 checkpoint 下的增量寫列表。delete_thread / delete_for_runs:320-348— 執行緒級清理,帶 DeltaChannel 警告。WRITES_IDX_MAP:795—ERROR/SCHEDULED/INTERRUPT/RESUME這四個控制 channel 用負 idx 占坑,避免和真任務寫撞 primary key。
資料流
PregelLoop 每跑完一個超步,會呼叫 saver 把當前所有通道值和這步新版本號寫下去。這裡看 put 的預設簽名和 InMemorySaver 的實作就夠清楚——引擎呼叫 put,saver 把 channel_values 按 (thread_id, ns, channel, version) 拆出去存 blob,主體進主表,回傳新 config 指向剛落的 checkpoint_id。
def put(
self,
config: RunnableConfig,
checkpoint: Checkpoint,
metadata: CheckpointMetadata,
new_versions: ChannelVersions,
) -> RunnableConfig:
"""Store a checkpoint with its configuration and metadata."""
raise NotImplementedError這是基類裡的純介面 (put method:277-298),不帶預設實作。每個子類別得自己決定哪些欄位進主表、哪些進 blob、用不用 ON CONFLICT DO UPDATE 之類。讀回來時 get_tuple 要負責把 blob 反序列化拼回 channel_values,並把同一 checkpoint 下的 pending_writes 一起裝進 CheckpointTuple (get_tuple:239-251)。
邊界與失敗
put_writes用負 idx 占坑控制 channel:WRITES_IDX_MAP把ERROR/SCHEDULED/INTERRUPT/RESUME映射到 -1 / -2 / -3 / -4 (WRITES_IDX_MAP:795),自訂 saver 的主鍵(task_id, idx)必須能接受負值,否則控制寫會撞主鍵衝突。get預設呼叫get_tuple:子類別只實作get_tuple就夠了,get預設實作幫你取.checkpoint(get:227-237)。但get_tuple不實作就是raise NotImplementedError,不能 silently 回傳None。- async 預設是
NotImplementedError:同步 saver 不實作aput時,預設 async 路徑不會自動轉執行緒——要嘛子類別自己寫asyncio.to_thread轉發,要嘛呼叫方呼叫同步路徑 (async defaults:468-509)。 metadata不能丟:list的filter參數就是按 metadata 欄位過濾的 (list method:253-275),saver 實作如果不存或存錯 metadata,過濾和分頁都會啞掉。- DeltaChannel 下不能 naive 刪:基類
prune的 docstring 明確說 (prune:374-415)——只留最新 checkpoint 會讓 delta channel 重建靜默回傳空,不報錯。自訂 saver 實作刪除前要走祖先鏈。 checkpoint_ns預設空串:子圖呼叫會帶非空 ns,頂層呼叫是""。saver 的主鍵必須把 ns 算進去,不然不同子圖的同名 checkpoint 會互相覆蓋。
小結
BaseCheckpointSaver 把「檢查點怎麼存、怎麼讀、怎麼列、怎麼寫增量」這幾件事一次性釘成五方法 + 一份 TypedDict 形狀,Pregel 引擎只跟介面打交道,後端實作任意換。具體兩種參考實作見兄弟頁面 InMemorySaver + SqliteSaver 和 PostgresSaver;執行時怎麼決定落盤節奏見 Pregel 引擎 和 PregelLoop。
對照官方資料:LangGraph 文件 · README