Skip to content

InMemorySaver + SqliteSaver:プロセス内とファイルレベルのチェックポイント

源码版本1.2.9

役割

BaseCheckpointSaver がインタフェース形状を与え、真に実装された最もよく使われる 2 つの saver が InMemorySaverSqliteSaver です。前者は状態全体をプロセスメモリ内の defaultdict に詰め込み、後者は同じフィールドを checkpointswrites の 2 枚の SQLite テーブルに展開します。両者で開発デバッグと単一マシン永続化の 2 シナリオをカバーします。プロダクション規模になると PostgresSaver に切り替えます。

InMemorySaverlibs/checkpoint/langgraph/checkpoint/memory/__init__.py にあります (InMemorySaver:33)。内部は 3 つの defaultdict で 3 種類のものを保持します:storage は checkpoint 本体 + メタデータ + parent_id、writes は (thread_id, ns, checkpoint_id) → (task_id, idx) → 書き込み記録、blobs は (thread_id, ns, channel, version) → シリアライズ後のバイナリ (storage / writes / blobs:68-83)。この 3 バケツ構造は Postgres/SQLite の 3 テーブル構造と 1 対 1 で対応し、表を dict で代用したものです。

SqliteSaverlibs/checkpoint-sqlite/langgraph/checkpoint/sqlite/__init__.py にあります (SqliteSaver:45)。BaseCheckpointSaver[str] を直接継承し、sqlite3.Connection を受け取り、setup() で 2 枚のテーブルを作ります(setup:129-166)。put / put_writes は標準 SQL INSERT OR REPLACE / INSERT OR IGNORE で走ります。threading.Lock で単一接続を直列化するため、check_same_thread=False でも安全です。

また MemorySaver = InMemorySaver は歴史的エイリアスです (MemorySaver alias:631)。古いコードでよく見る from langgraph.checkpoint.memory import MemorySaver が取得するのは同じクラスです。

設計動機

  • InMemory の 3 バケツ dict は 3 テーブルに対応:storage / writes / blobs の 3 フィールド (data members:68-83) は SQLite の checkpoints / writes 表および Postgres の checkpoints / checkpoint_writes / checkpoint_blobs 3 テーブルと 1 対 1 で対応します。InMemory が blob を独立して保存するのは、channel value のシリアライズをバージョン変化時に 1 回だけ行うためです。毎回の put で全チャネルを再シリアライズする事態を避けます。
  • PersistentDict でオプションのディスク書き込み:memory モジュールには PersistentDict(defaultdict) (PersistentDict:634) があり、pickle で dict をファイルにダンプし、再起動時に load します。InMemorySaver__init__factory パラメータ (__init__:85-99) を受け取るのは、factory=PersistentDict のとき自動的にコンテキストマネージャを装着するためです。これは「デバッグ用途だが再起動をまたいで状態を保ちたい」という折衷案です。
  • SQLite は WAL + 単一 lock:setup() の 1 行目は PRAGMA journal_mode=WAL (WAL:141) で、読みが書きをブロックしないようにします。ただし書きは self.lock で直列化します。sqlite3.Connection はスレッドセーフではないためです。docstring は「マルチスレッドにスケールしない」(note:48-54) と明記し、並行性が必要なら Postgres を使います。
  • 書き込みテーブルで REPLACE と IGNORE を使い分け:put_writes は SQLite で、書き込みがすべて WRITES_IDX_MAP に入るかどうかで INSERT OR REPLACEINSERT OR IGNORE を選びます (put_writes query:462-466)。制御 channel の書き込み(ERROR / INTERRUPT / RESUME)は上書き可能、通常の業務書き込みは衝突したら無視します。语义は InMemory の if inner_key[1] >= 0 ... continue と一致します (InMemorySaver.put_writes:499-509)。
  • InMemory の max(checkpoints.keys()) を latest に:明示的な ORDER BY はなく、checkpoint_id の文字列単調増加に頼ります (latest:282-283)。したがって checkpoint_id は辞書式順序で比較可能な UUID6/UUID7 でなければならず、純粋なランダム UUID は使えません。

主要ファイル

データフロー

InMemorySaver.put は「channel_values を blobs に分ける」と「本体を storage に入れる」の 2 つを一度に済ませます。saver が Checkpoint TypedDict を 3 層構造にどう分解するかが明晰に見えます:

python
c = checkpoint.copy()
thread_id = config["configurable"]["thread_id"]
checkpoint_ns = config["configurable"]["checkpoint_ns"]
values: dict[str, Any] = c.pop("channel_values")  # type: ignore[misc]
for k, v in new_versions.items():
    self.blobs[(thread_id, checkpoint_ns, k, v)] = (
        self.serde.dumps_typed(values[k]) if k in values else ("empty", b"")
    )
self.storage[thread_id][checkpoint_ns].update(
    {
        checkpoint["id"]: (
            self.serde.dumps_typed(c),
            self.serde.dumps_typed(get_checkpoint_metadata(config, metadata)),
            config["configurable"].get("checkpoint_id"),  # parent
        )
    }
)

(put body:448-464) blobs(thread_id, checkpoint_ns, channel, version) をキーにし、storage(thread_id, checkpoint_ns, checkpoint_id) をキーにします。2 つの索引は重ならず、channel value は checkpoint 本体と切り離され、バージョンが変わっていないチャネルは再シリアライズ不要です。読み戻し時、get_tuple_load_blobs を呼んで channel_versions 内の各 (channel, version) に対応する blob をデシリアライズして channel_values に組み戻し (_load_blobs:259-265)、CheckpointTuple に詰めます。

SqliteSaver は同じ論理で走りますが、dict が SQL に代わり、blob バケツは主表の checkpoint BLOB 列 + writes 表の value BLOB 列に代わります。独立した blob 表はなく、すべての channel value は直接 serde されて checkpoint blob フィールドに入ります:

python
cur.execute(
    "INSERT OR REPLACE INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, parent_checkpoint_id, type, checkpoint, metadata) VALUES (?, ?, ?, ?, ?, ?, ?)",
    (
        str(config["configurable"]["thread_id"]),
        checkpoint_ns,
        checkpoint["id"],
        config["configurable"].get("checkpoint_id"),
        type_,
        serialized_checkpoint,
        serialized_metadata,
    ),
)

(SqliteSaver.put INSERT:425-436)

境界と失敗

  • InMemorySaver はプロセスが落ちたら全消失:docstring は「InMemorySaver はデバッグやテストにのみ使う」と明記します (warning:40-44)。永続化するなら Sqlite か Postgres に替えます。
  • InMemorySaver の latest 取得は checkpoint_id の辞書式順序に依存:max(checkpoints.keys()) (max:282-283) は ID が単調増加すると仮定します。カスタム saver が非単調 ID を使うと、get_tuple が checkpoint_id なしで呼ばれたときに行を取り間違えます。
  • SqliteSaver は単一接続・単一ロック:cursor() コンテキストマネージャは self.lock を取得し (cursor:181-189)、並行書き込みは直列化され、長いトランザクションは他の読みをブロックします。並行性が必要なら Postgres の pipeline モードを使います。
  • SqliteSaver は ? プレースホルダで列名順序を硬编码:INSERT INTO checkpoints (...) の列リストと順序を一致させる必要があります (INSERT columns:426)。マイグレーションで列を追加するときは両方を更新します。task_path 列は後から追加されたもので、Postgres MIGRATIONS の末尾を参照してください (task_path migration:90)。
  • AsyncSqliteSaver は単なる to_thread ラップではない:AsyncSqliteSaver:38 は独立クラスで、aput / aget_tuple / alist はすべて aiosqlite で書き直されます (aput:509)。基底クラスのデフォルトである asyncio.to_thread 転送ではないため、非同期パスは真に並行し、グローバル lock で直列化されません。
  • factory= で PersistentDict を使うときは終了時に close() が必要:PersistentDict.sync() は pickle をファイルにダンプし (sync:657)、__exit__ を呼ばないと自動でディスクに落ちません。with InMemorySaver(factory=PersistentDict, ...) が慣用句です。

まとめ

InMemorySaver は「3 バケツ dict」の参考実装で、SqliteSaver は同じ 3 バケツを SQLite の 2 テーブル + WAL に展開します。両者でデバッグと単一マシン永続化をカバーします。主キー設計、blob 分割、制御 channel の負 idx といった取り決めは両実装で見えてきます。プロダクション規模は PostgresSaver を、インタフェース形状は BaseCheckpointSaver を参照してください。

公式資料:LangGraph 文档 · README