PostgresSaver : checkpoint relationnel de niveau production
Responsabilités
PostgresSaver est l'implémentation de BaseCheckpointSaver sur Postgres, et aussi le saver de production officiellement recommandé par LangGraph. Le code vit dans libs/checkpoint-postgres/langgraph/checkpoint/postgres/ ; l'entrée synchrone est PostgresSaver dans libs/checkpoint-postgres/langgraph/checkpoint/postgres/__init__.py (PostgresSaver:40), l'asynchrone est AsyncPostgresSaver dans aio.py (AsyncPostgresSaver:40), les deux partageant la classe de base BasePostgresSaver (BasePostgresSaver:297) qui détient les templates SQL et les scripts de migration.
Contrairement à l'approche SQLite « une table + dédup par clé primaire », PostgresSaver éclate les données en trois tables : checkpoints (table principale, JSONB), checkpoint_blobs (valeurs binaires des canaux), checkpoint_writes (écritures incrémentales par tâche) (MIGRATIONS:43-91). Cette répartition permet à la lecture de récupérer en une seule SELECT + sous-requêtes latérales l'intégralité « table principale + blobs + writes » (SELECT_SQL:93-118) et évite les N+1 de SQLite ; par ailleurs, la colonne blob est indépendante en BYTEA et ne participe pas à la désérialisation JSONB, donc un gros objet ne plombe pas la lecture de la table principale.
PostgresSaver supporte en plus le mode Pipeline de psycopg3 (supports_pipeline:60) : plusieurs statements sont envoyés en lot, le serveur ne renvoie pas les résultats immédiatement — adapté aux combinaisons executemany(blob) + execute(checkpoint) à l'intérieur d'un même put.
Motivation de conception
- Trois tables séparées, JSONB pour le corps : la table principale
checkpointsutilise une colonne JSONB pour le TypedDictCheckpoint, etcheckpoint_blobsstocke les gros objets indépendamment (checkpoint_blobs:57-65). Cela évite de fourrer toutchannel_valuessérialisé dans le JSONB — approche qui s'effondre en performance quand les canaux sont nombreux et les valeurs volumineuses. - Une seule SELECT ramène tout le tuple :
SELECT_SQLutilise deux sous-requêtes latérales (array_aggsurcheckpoint_blobs+array_aggsurcheckpoint_writes) pour récupérer toutes les données en un seul passage (lateral subqueries:101-117) ; côté Python,_load_checkpoint_tuplese contente de déballer les arrays, sans seconde requête. - UPSERT distingue
ON CONFLICT DO UPDATE/NOTHINGselon la sémantique :UPSERT_CHECKPOINT_BLOBS_SQLutiliseDO NOTHING(UPSERT_BLOBS:131-135),UPSERT_CHECKPOINTS_SQLutiliseDO UPDATE(UPSERT_CHECKPOINTS:137-144),UPSERT_CHECKPOINT_WRITES_SQLutiliseDO UPDATE, tandis queINSERT_CHECKPOINT_WRITES_SQLutiliseDO NOTHING(INSERT_WRITES:155-159). Les écritures sur canaux de contrôle doivent écraser l'état le plus récent (UPDATE) ; les écritures métier en collision sont ignorées (NOTHING) — sémantique alignée sur leif inner_key[1] >= 0 ... continued'InMemory. - Liste de migrations par numéro de version :
MIGRATIONSest une list (MIGRATIONS:43-91) ;setup()lit la versionvmaximale de la tablecheckpoint_migrationset n'exécute que les migrations suivantes (setup:85-110). Ajout de colonne, d'index ou de contrainte se fait en appendant une nouvelle chaîne de migration, sans toucher aux anciennes. - Le mode Pipeline offre de la vraie concurrence sur une connexion : en mode synchrone,
PostgresSaver._cursorprend le contextconn.pipeline()quandpipeline=True(pipeline branch:424-432) : plusieurs cursors sur la même connexion ne se bloquent pas mutuellement ; et le modepipe(issu defrom_conn_string(pipeline=True)) permet à une connexion d'être réutilisée par plusieurs threads (pipe branch:413-423), couplé àself.lockqui sérialise le cursor mais pas les statements.
Fichiers clés
MIGRATIONS:43-91—CREATEdes trois tables + migrations WAL/index/task_path.SELECT_SQL:93-118— requête principale, sous-requêtes latérales quiarray_aggblobs et writes ensemble.SELECT_PENDING_SENDS_SQL:120-129— récupère les writes PUSH en attente, triés par task_path / task_id / idx.UPSERT / INSERT SQL:131-159— quatre statements, distinguantDO UPDATEetDO NOTHING.BasePostgresSaver:297— classe de base partagée, détient les constantes de templates SQL.PostgresSaver 类:40-60— implémentation synchrone, accepteConnectionouConnectionPool, embarque unthreading.Lock.from_conn_string:62-83— factory en context manager,pipeline=Trueactiveconn.pipeline().setup:85-110— lit la version de migration, exécute dans l'ordre celles non appliquées et enregistre leur v.put:263-345— déplace_DeltaSnapshot/ valeurs non primitives vers blob_values, met le corps dans le JSONB._cursor:405-442— choisit le context selon le mode pipeline / pipe / normal.AsyncPostgresSaver:40— variante async,aput/aget_tuple/alisten async natif.
Flux de données
PostgresSaver.put est un bon exemple pour voir comment un saver éclate un TypedDict Checkpoint : il décide d'abord quelles valeurs de canal sont inline dans le JSONB principal, lesquelles sont déplacées vers checkpoint_blobs, puis à l'intérieur d'un curseur exécute 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) Notez le chemin spécial pour _DeltaSnapshot : la valeur du canal n'a dans la table principale qu'un placeholder True, les données réelles vont dans blob — c'est la base sur laquelle DeltaChannel remonte la chaîne d'ancêtres lors de la reconstruction. À la lecture, get_tuple exécute SELECT_SQL ; psycopg convertit le JSONB en dict et le BYTEA en bytes automatiquement, puis côté Python _load_checkpoint_tuple recompose le CheckpointTuple (_load_checkpoint_tuple:552).
Limites et échecs
from_conn_string(pipeline=True)ne se mélange pas avecConnectionPool:__init__lève explicitement (pool vs pipe check:52-55) : pipeline est un mode connexion unique, pool est déjà multi-connexions.setup()doit être appelé explicitement au moins une fois : le docstring précise « MUST be called directly by the user the first time checkpointer is used » (setup docstring:85-91). Il crée les tables et joue les migrations ; la base doit avoir les permissions adéquates au préalable.- L'ordre de
MIGRATIONSne doit pas changer : les migrations utilisent la position dans la list comme numéro de version (MIGRATIONS:43-91) ; une chaîne déjà publiée ne peut plus être modifiée, sinon levde la tablecheckpoint_migrationsse décale par rapport à la position dans le code. CREATE INDEX CONCURRENTLYéchoue dans une transaction : les deux statements d'index concurrent dansMIGRATIONS(CONCURRENTLY indexes:82-89) s'exécutent dans le même contextesetup(); si la connexion n'est pas en autocommit, erreur. C'est pour cela quefrom_conn_stringforceautocommit=True(autocommit:76-78).supports_pipelinedétecté à l'exécution :Capabilities().has_pipeline()(has_pipeline:60) est appelé dans__init__pour vérifier la version de psycopg et les capacités côté serveur ; en cas de psycopg ancien ou de certains Postgres managés, dégradation vers le mode transaction normal (transaction fallback:434-439).- DeltaChannel a une fast path dédiée :
get_delta_channel_historyest surchargée dansPostgresSaver(get_delta_channel_history:444-463) en une requête en deux phases (Stage 1 : pagination metadata + Stage 2 : UNION ALL pour ramener les writes) ; elle ne suit pas l'implémentation par défaut de la classe de base. Un saver personnalisé doit implémenter cette interface synchrone pour pouvoir faire tourner DeltaChannel.
Résumé
PostgresSaver est le saver de production recommandé par LangGraph ; la séparation en trois tables permet à « JSONB principal + BYTEA pour gros objets + writes incrémentales » d'aller chacun vers le stockage le plus adapté. Le mode Pipeline et l'implémentation async native offrent le débit attendu à l'échelle production. Pour la forme de l'interface, voir BaseCheckpointSaver ; pour les scénarios légers, voir InMemorySaver + SqliteSaver.
Voir la documentation officielle : documentation LangGraph · README