Skip to content

PostgresSaver:プロダクション級リレーショナルチェックポイント

源码版本1.2.9

役割

PostgresSaverBaseCheckpointSaver の Postgres 向け実装であり、LangGraph が公式に推奨するプロダクション用 saver です。コードは libs/checkpoint-postgres/langgraph/checkpoint/postgres/ にあり、同期版のエントリは libs/checkpoint-postgres/langgraph/checkpoint/postgres/__init__.pyPostgresSaver (PostgresSaver:40)、非同期版は aio.pyAsyncPostgresSaver (AsyncPostgresSaver:40) で、共通の基底クラス BasePostgresSaver (BasePostgresSaver:297) が SQL テンプレートとマイグレーションスクリプトを担います。

SQLite の単一テーブル + 主キー重複除外というやり方とは異なり、PostgresSaver はデータを 3 枚のテーブルに分割します:checkpoints(主表、JSONB)、checkpoint_blobs(channel value のバイナリ)、checkpoint_writes(task の差分書き込み)(MIGRATIONS:43-91)。この分割により、読み出しは 1 本の SELECT + lateral サブクエリで主表 + blobs + writes を一度に取り戻せます (SELECT_SQL:93-118)。SQLite のような N+1 クエリを避け、同時に blob 列は BYTEA として独立し JSONB のデシリアライズに加わらないため、大きなオブジェクトが主表の読み出しを引きずりません。

PostgresSaver はさらに psycopg3 の Pipeline モードをサポートします (supports_pipeline:60)。複数文をまとめて送り、サーバ側は即座に結果を返さないモードで、1 回の put 内部で executemany(blob) + execute(checkpoint) という組み合わせを流すのに適しています。

設計動機

  • 3 テーブル分離、JSONB で本体を保存:checkpoints 主表は JSONB 列で Checkpoint TypedDict を格納し、checkpoint_blobs が大きなオブジェクトを独立して保持します (checkpoint_blobs:57-65)。これにより channel_values 全体を JSONB に詰め込むことを避けられます。そのようなやり方はチャネル (channel) が多く value が大きいとき性能が崩れます。
  • 1 本の SELECT で tuple を完成:SELECT_SQL は 2 つの lateral サブクエリ(array_agg over checkpoint_blobs + array_agg over checkpoint_writes)で全データを一度に取り戻します (lateral subqueries:101-117)。_load_checkpoint_tuple は Python 側で array を展開するだけで、二次クエリは不要です。
  • UPSERT を ON CONFLICT DO UPDATE/NOTHING で语义分け:UPSERT_CHECKPOINT_BLOBS_SQLDO NOTHING (UPSERT_BLOBS:131-135)、UPSERT_CHECKPOINTS_SQLDO UPDATE (UPSERT_CHECKPOINTS:137-144)、UPSERT_CHECKPOINT_WRITES_SQLDO UPDATE、そして INSERT_CHECKPOINT_WRITES_SQLDO NOTHING (INSERT_WRITES:155-159) を使います。制御 channel の書き込みは最新状態で上書き(UPDATE)、通常の業務書き込みは衝突したら無視(NOTHING)という语义で、InMemory の if inner_key[1] >= 0 ... continue と揃えます。
  • マイグレーション一覧はバージョン番号を追加:MIGRATIONS は list であり (MIGRATIONS:43-91)、setup()checkpoint_migrations 表の最大 v を読み、それ以降のマイグレーションのみ実行します (setup:85-110)。列追加、インデックス追加、制約変更はすべて migration 文字列を追加する形で行い、古いマイグレーションは書き換えません。
  • Pipeline モードで単接続の真の並行:同期 PostgresSaver._cursorpipeline=True のとき conn.pipeline() コンテキストを取得し (pipeline branch:424-432)、1 接続上の複数 cursor が互いにブロックしません。一方 pipe モード(from_conn_string(pipeline=True) から)は 1 接続をマルチスレッドで再利用し (pipe branch:413-423)、self.lock で cursor を直列化しつつ文は直列化しません。

主要ファイル

  • MIGRATIONS:43-91 — 3 テーブルの CREATE 文 + WAL/インデックス/task_path マイグレーション。
  • SELECT_SQL:93-118 — 主クエリ。lateral サブクエリで blobs と writes を一緒に array_agg します。
  • SELECT_PENDING_SENDS_SQL:120-129 — 送信待ち PUSH タスク writes を取得。task_path / task_id / idx 順。
  • UPSERT / INSERT SQL:131-159 — 4 本の書き込み文。DO UPDATE / DO NOTHING を使い分けます。
  • BasePostgresSaver:297 — 共通基底クラス。SQL テンプレート定数を保持します。
  • PostgresSaver 类:40-60 — 同期実装。Connection または ConnectionPool を受け取り、threading.Lock を持ちます。
  • from_conn_string:62-83 — コンテキストマネージャファクトリ。pipeline=Trueconn.pipeline() を使います。
  • setup:85-110 — マイグレーションバージョンを読み、未適用のマイグレーションを順に実行して v を記録します。
  • put:263-345_DeltaSnapshot や非原始値を blob_values に移し、本体は JSONB に入れます。
  • _cursor:405-442 — pipeline / pipe の有無で 3 モードのコンテキストを使い分けます。
  • AsyncPostgresSaver:40 — 非同期版。aput / aget_tuple / alist がネイティブ async です。

データフロー

PostgresSaver.put は saver が Checkpoint TypedDict をどう分割するかを見る好例です。どの channel value を JSONB 本体にインライン化し、どれを checkpoint_blobs に移すかを決め、1 つの cursor 内で executemany(blob) + execute(checkpoint) を走らせます:

python
# inline primitive values in checkpoint table
# others are stored in blobs table
blob_values = {}
for k, v in checkpoint["channel_values"].items():
    if isinstance(v, _DeltaSnapshot):
        blob_values[k] = copy["channel_values"].pop(k)
        copy["channel_values"][k] = True
    elif v is None or isinstance(v, (str, int, float, bool)):
        pass
    else:
        blob_values[k] = copy["channel_values"].pop(k)

with self._cursor(pipeline=True) as cur:
    if blob_versions := {
        k: v for k, v in new_versions.items() if k in blob_values
    }:
        cur.executemany(
            self.UPSERT_CHECKPOINT_BLOBS_SQL,
            self._dump_blobs(
                thread_id,
                checkpoint_ns,
                blob_values,
                blob_versions,
            ),
        )
    cur.execute(
        self.UPSERT_CHECKPOINTS_SQL,
        (
            thread_id,
            checkpoint_ns,
            checkpoint["id"],
            checkpoint_id,
            Jsonb(copy),
            Jsonb(get_serializable_checkpoint_metadata(config, metadata)),
        ),
    )

(put body:309-344) _DeltaSnapshot は特殊パスを通ります。チャネル値は主表に True のプレースホルダだけを置き、実データは blob に入ります。これが DeltaChannel が再構築時に祖先チェーンを逆走する根拠です。読み戻し時には get_tupleSELECT_SQL を走らせ、psycopg が JSONB を自動で dict に、BYTEA を bytes に変換し、Python 側で _load_checkpoint_tupleCheckpointTuple を組み立てます (_load_checkpoint_tuple:552)。

境界と失敗

  • from_conn_string(pipeline=True) は ConnectionPool と併用不可:__init__ は明示的に raise します (pool vs pipe check:52-55)。pipeline は単接続モードで、pool はすでに多接続です。
  • setup() は明示的に 1 回呼ぶ必要あり:docstring に「MUST be called directly by the user the first time checkpointer is used」と書かれています (setup docstring:85-91)。テーブル作成 + マイグレーション実行を行うので、事前にデータベースへ権限を与えておく必要があります。
  • MIGRATIONS の順序変更は破綻する:マイグレーションは list の位置をバージョン番号とします (MIGRATIONS:43-91)。一度公開したマイグレーション文字列は変更できません。変更すると checkpoint_migrations 表の v とコード上の位置がずれます。
  • CREATE INDEX CONCURRENTLY はトランザクション内で失敗する:MIGRATIONS の 2 本の並行インデックス作成文 (CONCURRENTLY indexes:82-89) は他のマイグレーションと同じ setup() コンテキストで走るため、接続が autocommit モードでないとエラーになります。from_conn_string が明示的に autocommit=True にするのはこのためです (autocommit:76-78)。
  • supports_pipeline は実行時検出:Capabilities().has_pipeline() (has_pipeline:60) は __init__ で psycopg のバージョンとサーバ側能力を調べ、古い psycopg や一部のマネージド Postgres では通常トランザクションモードに退化します (transaction fallback:434-439)。
  • DeltaChannel 専用の高速パス:get_delta_channel_historyPostgresSaver でオーバーライドされ (get_delta_channel_history:444-463)、2 段階クエリ(Stage 1 metadata ページング + Stage 2 UNION ALL で writes を取得)で走ります。基底クラスのデフォルト実装を使わないため、カスタム saver はこのインタフェースを実装しないと DeltaChannel を動かせません。

まとめ

PostgresSaver は LangGraph 推奨のプロダクション saver です。3 テーブル分割で「主表 JSONB + 大对象 BYTEA + 差分 writes」がそれぞれ最適なストレージに流れます。Pipeline モードとネイティブ async 実装がプロダクション規模のスループットを支えます。インタフェースの形状は BaseCheckpointSaver を、軽量ケースは InMemorySaver + SqliteSaver を参照してください。

公式資料:LangGraph 文档 · README