PostgresSaver: checkpoint relacional de producción
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
checkpointsusa una columna JSONB para el TypedDictCheckpoint;checkpoint_blobsguarda los objetos grandes de forma independiente (checkpoint_blobs:57-65). Así se evita meter todochannel_valuesserializado dentro del JSONB —algo que colapsa en rendimiento cuando hay muchos canales o valores grandes. - Un único SELECT trae el tuple completo:
SELECT_SQLusa dos subconsultas laterales (array_aggsobrecheckpoint_blobs+array_aggsobrecheckpoint_writes) para obtener todos los datos en una sola consulta (lateral subqueries:101-117);_load_checkpoint_tuplesólo desempaqueta los arrays en Python, sin segundas consultas. - UPSERT con
ON CONFLICT DO UPDATE/NOTHINGpara distinguir semánticas:UPSERT_CHECKPOINT_BLOBS_SQLusaDO NOTHING(UPSERT_BLOBS:131-135),UPSERT_CHECKPOINTS_SQLusaDO UPDATE(UPSERT_CHECKPOINTS:137-144),UPSERT_CHECKPOINT_WRITES_SQLusaDO UPDATE, yINSERT_CHECKPOINT_WRITES_SQLusaDO 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 elif inner_key[1] >= 0 ... continuede InMemory. - Lista de migraciones con versión incremental:
MIGRATIONSes una list (MIGRATIONS:43-91);setup()lee el v máximo de la tablacheckpoint_migrationsy 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:
_cursordelPostgresSaversíncrono, conpipeline=True, entra en el contextoconn.pipeline()(pipeline branch:424-432) —varios cursores sobre una conexión no se bloquean entre sí; el modopipe(desdefrom_conn_string(pipeline=True)) permite reutilizar una conexión entre varios hilos (pipe branch:413-423), conself.lockserializando 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 conarray_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 distinguenDO UPDATE/DO NOTHING.BasePostgresSaver:297— clase base compartida; contiene las constantes con las plantillas SQL.PostgresSaver 类:40-60— implementación síncrona; recibeConnectionoConnectionPooly traethreading.Lock.from_conn_string:62-83— fábrica como gestor de contexto; conpipeline=Trueusaconn.pipeline().setup:85-110— lee la versión de migraciones, aplica las pendientes en orden y registra el v.put:263-345— mueve_DeltaSnapshoty valores no primitivos ablob_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/alistson 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):
# 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 encheckpoint_migrationsse desincronizaría con la posición en el código. CREATE INDEX CONCURRENTLYfalla dentro de una transacción: dos sentencias de creación concurrente de índices enMIGRATIONS(CONCURRENTLY indexes:82-89) corren dentro del mismo contextosetup(); si la conexión no está en autocommit fallan.from_conn_stringfuerzaautocommit=Truepor esto mismo (autocommit:76-78).supports_pipelinese 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_historyse reescribe enPostgresSaver(get_delta_channel_history:444-463) con una consulta en dos fases (paginación de metadata en Stage 1 +UNION ALLde 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