Skip to content

BaseCheckpointSaver:永続化インタフェースの形状

源码版本1.2.9

役割

LangGraph のランタイムでは、Pregel エンジンは 1 つのスーパーステップ (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 は「このタスクがこのステップで生んだ書き込み」をチェックポイントとは別に格納します。書き込みは差分でディスクに落ちるため、毎回 snapshot 全体を書き直す必要はありません。

言い換えると、BaseCheckpointSaver は PregelLoop と具象ストレージバックエンドの間の契約です。エンジンはこのインタフェースだけを知っていれば、バックエンドの実装を替えてもエンジン側は変更不要です。バックエンド側は 5 つのメソッドを埋めるだけで、そのバックエンドの並行性と永続化の特性をエンジンが自動的に獲得します。

設計動機

  • 状態と書き込みを分離:put は snapshot 全体を、put_writes は 1 つの task があるステップで生んだ差分書き込みを格納します (put_writes:300-318)。分けた理由は、中断恢復時に「このステップでどの task がすでに書いたか、どれがまだか」を精確に復元するためです。前のチェックポイントへ全体を巻き戻すしかない、という事態を避けます。
  • thread_id を主キーに:put が返す config には thread_id / checkpoint_ns / checkpoint_id だけが入ります (put:277-298)。呼び出しをまたいでも、プロセスをまたいでも、この 3 フィールドで 1 つの checkpoint を位置づけます。checkpoint_ns はサブグラフ (subgraph) 向けの名前空間で、空文字列ならトップレベルグラフです。
  • 同期 + 非同期の二系統インタフェース:各書き込みメソッドには aput / aget / alist / aput_writes の非同期版があります (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 再構築が暗黙に空を返します。

主要ファイル

  • BaseCheckpointSaver 类:176 — 抽象基底クラス本体。put/get/list/put_writes などのインタフェースを規定します。
  • Checkpoint TypedDict:92-123 — 1 つの checkpoint のフィールド形状:channel_values/channel_versions/versions_seen/id/ts。
  • CheckpointTuple:139-146get_tuple の戻り値。checkpoint、metadata、parent_config、pending_writes を 1 つの NamedTuple に詰めます。
  • CheckpointMetadata:38-86source / step / parents / run_id / DeltaChannel 用の counters_since_delta_snapshot
  • get 方法:227-237get はデフォルトで get_tuple に委譲し、checkpoint フィールドだけを返します。
  • put 方法:277-298 — checkpoint + metadata 全体を保存し、更新後の config を返します。
  • put_writes 方法:300-318 — ある checkpoint における task の差分書き込みリストを保存します。
  • delete_thread / delete_for_runs:320-348 — スレッド単位のクリーンアップ。DeltaChannel の警告付き。
  • WRITES_IDX_MAP:795ERROR / SCHEDULED / INTERRUPT / RESUME の 4 つの制御 channel は負の idx で枠を取り、本物の task 書き込みと主キー衝突しないようにします。

データフロー

PregelLoop が 1 つのスーパーステップを走ら終えると、saver を呼んで現在のすべてのチャネル値とこのステップの新バージョン番号を書き込みます。ここでは put のデフォルトシグネチャと InMemorySaver の実装を見れば明晰です。エンジンは put を呼び、saver は channel_values を (thread_id, ns, channel, version) ごとに分けて blob として保存し、本体は主表に入れ、今落ちた checkpoint_id を指す新 config を返します。

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 で、暗黙に None を返すことはありません。
  • async のデフォルトは NotImplementedError:同期 saver が aput を実装しない場合、デフォルトの async パスは自動的にスレッド転送しません。サブクラスが自分で asyncio.to_thread で転送するか、呼び出し側が同期パスを呼ぶ必要があります (async defaults:468-509)。
  • metadata は失えない:listfilter パラメータは metadata フィールドでフィルタします (list method:253-275)。saver が metadata を保存しない、あるいは保存を間違えると、フィルタとページングが黙って効かなくなります。
  • DeltaChannel 配下では naively 削除禁止:基底クラスの prune docstring は明示的に言及します (prune:374-415)。最新 checkpoint だけ残すと delta channel の再構築が暗黙に空を返し、エラーも出ません。カスタム saver は削除前に祖先チェーンを歩く必要があります。
  • checkpoint_ns はデフォルト空文字列:サブグラフ呼び出しは非空 ns を持ち、トップレベル呼び出しは "" です。saver の主キーは ns を含まなければならず、さもなくば異なるサブグラフの同名 checkpoint が互いに上書きされます。

まとめ

BaseCheckpointSaver は「チェックポイントをどう保存し、どう読み、どう列挙し、どう差分を書くか」という事柄を 5 つのメソッド + 1 つの TypedDict 形状にひとまとめにします。Pregel エンジンはインタフェースとだけ付き合い、バックエンド実装は自由に替えられます。2 つの参考実装は兄弟ページ InMemorySaver + SqliteSaverPostgresSaver を参照してください。ランタイムがどうディスク書き込みのペースを決めるかは Pregel エンジンPregelLoop を参照してください。

公式資料:LangGraph 文档 · README