PostgresSaver:プロダクション級リレーショナルチェックポイント
役割
PostgresSaver は BaseCheckpointSaver の Postgres 向け実装であり、LangGraph が公式に推奨するプロダクション用 saver です。コードは libs/checkpoint-postgres/langgraph/checkpoint/postgres/ にあり、同期版のエントリは libs/checkpoint-postgres/langgraph/checkpoint/postgres/__init__.py の PostgresSaver (PostgresSaver:40)、非同期版は aio.py の AsyncPostgresSaver (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 列でCheckpointTypedDict を格納し、checkpoint_blobsが大きなオブジェクトを独立して保持します (checkpoint_blobs:57-65)。これにより channel_values 全体を JSONB に詰め込むことを避けられます。そのようなやり方はチャネル (channel) が多く value が大きいとき性能が崩れます。 - 1 本の SELECT で tuple を完成:
SELECT_SQLは 2 つの lateral サブクエリ(array_aggovercheckpoint_blobs+array_aggovercheckpoint_writes)で全データを一度に取り戻します (lateral subqueries:101-117)。_load_checkpoint_tupleは Python 側で array を展開するだけで、二次クエリは不要です。 - UPSERT を
ON CONFLICT DO UPDATE/NOTHINGで语义分け:UPSERT_CHECKPOINT_BLOBS_SQLはDO NOTHING(UPSERT_BLOBS:131-135)、UPSERT_CHECKPOINTS_SQLはDO UPDATE(UPSERT_CHECKPOINTS:137-144)、UPSERT_CHECKPOINT_WRITES_SQLはDO UPDATE、そしてINSERT_CHECKPOINT_WRITES_SQLはDO 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._cursorはpipeline=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=Trueはconn.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) を走らせます:
# 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_tuple が SELECT_SQL を走らせ、psycopg が JSONB を自動で dict に、BYTEA を bytes に変換し、Python 側で _load_checkpoint_tuple が CheckpointTuple を組み立てます (_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_historyはPostgresSaverでオーバーライドされ (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