Skip to content

PostgresSaver : checkpoint relationnel de niveau production

源码版本1.2.9

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 checkpoints utilise une colonne JSONB pour le TypedDict Checkpoint, et checkpoint_blobs stocke les gros objets indépendamment (checkpoint_blobs:57-65). Cela évite de fourrer tout channel_values sé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_SQL utilise deux sous-requêtes latérales (array_agg sur checkpoint_blobs + array_agg sur checkpoint_writes) pour récupérer toutes les données en un seul passage (lateral subqueries:101-117) ; côté Python, _load_checkpoint_tuple se contente de déballer les arrays, sans seconde requête.
  • UPSERT distingue ON CONFLICT DO UPDATE/NOTHING selon la sémantique : UPSERT_CHECKPOINT_BLOBS_SQL utilise DO NOTHING (UPSERT_BLOBS:131-135), UPSERT_CHECKPOINTS_SQL utilise DO UPDATE (UPSERT_CHECKPOINTS:137-144), UPSERT_CHECKPOINT_WRITES_SQL utilise DO UPDATE, tandis que INSERT_CHECKPOINT_WRITES_SQL utilise DO 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 le if inner_key[1] >= 0 ... continue d'InMemory.
  • Liste de migrations par numéro de version : MIGRATIONS est une list (MIGRATIONS:43-91) ; setup() lit la version v maximale de la table checkpoint_migrations et 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._cursor prend le context conn.pipeline() quand pipeline=True (pipeline branch:424-432) : plusieurs cursors sur la même connexion ne se bloquent pas mutuellement ; et le mode pipe (issu de from_conn_string(pipeline=True)) permet à une connexion d'être réutilisée par plusieurs threads (pipe branch:413-423), couplé à self.lock qui sérialise le cursor mais pas les statements.

Fichiers clés

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) :

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) 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 avec ConnectionPool : __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 MIGRATIONS ne 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 le v de la table checkpoint_migrations se décale par rapport à la position dans le code.
  • CREATE INDEX CONCURRENTLY échoue dans une transaction : les deux statements d'index concurrent dans MIGRATIONS (CONCURRENTLY indexes:82-89) s'exécutent dans le même contexte setup() ; si la connexion n'est pas en autocommit, erreur. C'est pour cela que from_conn_string force autocommit=True (autocommit:76-78).
  • supports_pipeline dé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_history est surchargée dans PostgresSaver (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