Skip to content

PostgresSaver: checkpoint relacional de producción

源码版本1.2.9

Responsabilidades

PostgresSaver es la implementación de BaseCheckpointSaver sobre Postgres y el saver recomendado por LangGraph para producción. El código está en libs/checkpoint-postgres/langgraph/checkpoint/postgres/; la entrada síncrona es PostgresSaver en libs/checkpoint-postgres/langgraph/checkpoint/postgres/__init__.py (PostgresSaver:40), la asíncrona es AsyncPostgresSaver en aio.py (AsyncPostgresSaver:40) y comparten la base BasePostgresSaver (BasePostgresSaver:297), que contiene las plantillas SQL y los scripts de migración.

A diferencia del enfoque de SQLite —una sola tabla con deduplicación por primary key—, PostgresSaver descompone los datos en tres tablas: checkpoints (tabla principal, JSONB), checkpoint_blobs (valores binarios de canal) y checkpoint_writes (escrituras incrementales de tarea) (MIGRATIONS:43-91). Esta separación permite que la lectura use un único SELECT + subconsultas laterales para traer de una vez principal + blobs + writes (SELECT_SQL:93-118), evitando el patrón N+1 de SQLite; además, la columna de blob es BYTEA independiente y no participa en la deserialización JSONB, así que los objetos grandes no penalizan la lectura de la tabla principal.

PostgresSaver además soporta el modo Pipeline de psycopg3 (supports_pipeline:60) —varias sentencias se empaquetan y envían juntas sin esperar resultado inmediato del servidor, ideal para combinaciones como executemany(blob) + execute(checkpoint) dentro de un mismo put.

Motivación de diseño

  • Tres tablas separadas, JSONB para el cuerpo: la tabla principal checkpoints usa una columna JSONB para el TypedDict Checkpoint; checkpoint_blobs guarda los objetos grandes de forma independiente (checkpoint_blobs:57-65). Así se evita meter todo channel_values serializado dentro del JSONB —algo que colapsa en rendimiento cuando hay muchos canales o valores grandes.
  • Un único SELECT trae el tuple completo: SELECT_SQL usa dos subconsultas laterales (array_agg sobre checkpoint_blobs + array_agg sobre checkpoint_writes) para obtener todos los datos en una sola consulta (lateral subqueries:101-117); _load_checkpoint_tuple sólo desempaqueta los arrays en Python, sin segundas consultas.
  • UPSERT con ON CONFLICT DO UPDATE/NOTHING para distinguir semánticas: UPSERT_CHECKPOINT_BLOBS_SQL usa DO NOTHING (UPSERT_BLOBS:131-135), UPSERT_CHECKPOINTS_SQL usa DO UPDATE (UPSERT_CHECKPOINTS:137-144), UPSERT_CHECKPOINT_WRITES_SQL usa DO UPDATE, y INSERT_CHECKPOINT_WRITES_SQL usa DO NOTHING (INSERT_WRITES:155-159). Las escrituras a canales de control sobrescriben el estado más reciente (UPDATE); las escrituras de negocio normales se ignoran si colisionan (NOTHING) —alineado con el if inner_key[1] >= 0 ... continue de InMemory.
  • Lista de migraciones con versión incremental: MIGRATIONS es una list (MIGRATIONS:43-91); setup() lee el v máximo de la tabla checkpoint_migrations y ejecuta sólo las migraciones posteriores (setup:85-110). Añadir columnas, índices o restricciones se hace agregando cadenas de migración, sin tocar las antiguas.
  • Modo Pipeline para concurrencia real sobre una conexión: _cursor del PostgresSaver síncrono, con pipeline=True, entra en el contexto conn.pipeline() (pipeline branch:424-432) —varios cursores sobre una conexión no se bloquean entre sí; el modo pipe (desde from_conn_string(pipeline=True)) permite reutilizar una conexión entre varios hilos (pipe branch:413-423), con self.lock serializando cursores pero no sentencias.

Archivos clave

  • MIGRATIONS:43-91 — sentencias CREATE de las tres tablas + migraciones de WAL/índices/task_path.
  • SELECT_SQL:93-118 — consulta principal; las subconsultas laterales agregan blobs y writes con array_agg.
  • SELECT_PENDING_SENDS_SQL:120-129 — recupera las writes pendientes de PUSH, ordenadas por task_path / task_id / idx.
  • UPSERT / INSERT SQL:131-159 — cuatro sentencias de escritura que distinguen DO UPDATE / DO NOTHING.
  • BasePostgresSaver:297 — clase base compartida; contiene las constantes con las plantillas SQL.
  • PostgresSaver 类:40-60 — implementación síncrona; recibe Connection o ConnectionPool y trae threading.Lock.
  • from_conn_string:62-83 — fábrica como gestor de contexto; con pipeline=True usa conn.pipeline().
  • setup:85-110 — lee la versión de migraciones, aplica las pendientes en orden y registra el v.
  • put:263-345 — mueve _DeltaSnapshot y valores no primitivos a blob_values, y el cuerpo principal al JSONB.
  • _cursor:405-442 — selecciona el contexto entre los modos pipeline / pipe / normal.
  • AsyncPostgresSaver:40 — versión asíncrona; aput / aget_tuple / alist son async nativas.

Flujo de datos

PostgresSaver.put es un buen ejemplo para ver cómo un saver descompone el TypedDict Checkpoint: primero decide qué valores de canal se inlinean en el cuerpo JSONB y cuáles se mueven a checkpoint_blobs, y luego dentro de un mismo cursor ejecuta 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) Nota: _DeltaSnapshot sigue una ruta especial —en la tabla principal sólo se coloca un placeholder True y los datos reales van al blob; esa es la pista que usa DeltaChannel para recorrer la cadena de ancestros al reconstruir. Al leer de vuelta, get_tuple ejecuta SELECT_SQL: psycopg convierte el JSONB en dict y BYTEA en bytes automáticamente, y desde Python se llama a _load_checkpoint_tuple para ensamblar el CheckpointTuple (_load_checkpoint_tuple:552).

Límites y fallos

  • from_conn_string(pipeline=True) no se puede mezclar con ConnectionPool: __init__ lanza un error explícito (pool vs pipe check:52-55); pipeline es un modo de conexión única y pool ya es multi-conexión.
  • setup() debe invocarse una vez explícitamente: la docstring dice «MUST be called directly by the user the first time checkpointer is used» (setup docstring:85-91). Crea tablas y ejecuta migraciones; antes de correrlo hay que tener listos los permisos de base de datos.
  • Cambiar el orden de MIGRATIONS revienta: las migraciones usan la posición en la list como número de versión (MIGRATIONS:43-91); las cadenas de migración ya publicadas no se pueden modificar, o el v guardado en checkpoint_migrations se desincronizaría con la posición en el código.
  • CREATE INDEX CONCURRENTLY falla dentro de una transacción: dos sentencias de creación concurrente de índices en MIGRATIONS (CONCURRENTLY indexes:82-89) corren dentro del mismo contexto setup(); si la conexión no está en autocommit fallan. from_conn_string fuerza autocommit=True por esto mismo (autocommit:76-78).
  • supports_pipeline se detecta en runtime: Capabilities().has_pipeline() (has_pipeline:60) comprueba en __init__ la versión de psycopg y las capacidades del servidor; psycopg antiguo o ciertos Postgres gestionados degradan a modo de transacción normal (transaction fallback:434-439).
  • DeltaChannel tiene una vía rápida específica: get_delta_channel_history se reescribe en PostgresSaver (get_delta_channel_history:444-463) con una consulta en dos fases (paginación de metadata en Stage 1 + UNION ALL de writes en Stage 2), sin usar la implementación por defecto de la clase base. Un saver personalizado debe implementar también esta interfaz si quiere soportar DeltaChannel.

Resumen

PostgresSaver es el saver de producción recomendado por LangGraph. La separación en tres tablas permite que «JSONB principal + BYTEA para objetos grandes + writes incrementales» cada uno vaya al almacenamiento que mejor le sienta; el modo Pipeline y la implementación async nativa aportan el rendimiento que exige producción. Para la forma de la interfaz, véase BaseCheckpointSaver; para escenarios ligeros, InMemorySaver + SqliteSaver.

Véase la documentación oficial: documentación de LangGraph · README