Skip to content

BaseCheckpointSaver:持久化介面的形狀

源码版本1.2.9

職責

在 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_valueschannel_versionsversions_seenid / ts / v 這幾樣。put 收到的不只是一個 Checkpoint,還有 metadatanew_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 路徑走原生驅動。
  • 元資料帶來源標籤:CheckpointMetadatasource 欄位取 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 重建會靜默回傳空。

關鍵檔案

資料流

PregelLoop 每跑完一個超步,會呼叫 saver 把當前所有通道值和這步新版本號寫下去。這裡看 put 的預設簽名和 InMemorySaver 的實作就夠清楚——引擎呼叫 put,saver 把 channel_values 按 (thread_id, ns, channel, version) 拆出去存 blob,主體進主表,回傳新 config 指向剛落的 checkpoint_id。

python
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_MAPERROR / 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 不能丟:listfilter 參數就是按 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 + SqliteSaverPostgresSaver;執行時怎麼決定落盤節奏見 Pregel 引擎PregelLoop

對照官方資料:LangGraph 文件 · README