BaseCheckpointSaver: la forma de la interfaz de persistencia
Responsabilidades
En el runtime de LangGraph, el motor Pregel tras ejecutar cada superpaso (superstep) debe persistir los valores de todos sus canales (channel) junto con las escrituras (writes) producidas en ese paso; así, la próxima vez que se entre con el mismo thread_id se puede continuar desde el estado anterior, y también es posible recuperarse tras una interrupción (interrupt), reproducir el historial y hacer depuración con viaje en el tiempo. BaseCheckpointSaver es la interfaz abstracta de esa capa de «persistencia» (BaseCheckpointSaver:176). No le importa dónde guardas los datos —un dict en memoria, SQLite, Postgres, Redis, cualquiera sirve—; sólo estipula los métodos que un saver debe implementar: put / get / get_tuple / list / put_writes (put / get / list / put_writes:227-318).
Al mismo tiempo fija la forma del concepto «punto de control (checkpoint)» como el TypedDict Checkpoint: Checkpoint TypedDict:92-123, que contiene channel_values, channel_versions, versions_seen, id / ts / v. Lo que put recibe no es sólo un Checkpoint, sino también metadata y new_versions; el saver decide qué campos van a la tabla principal y cuáles a la tabla de blobs. put_writes por su parte guarda «las escrituras que produjo esta tarea en este paso», desacopladas del checkpoint, de modo que las escrituras se puedan persistir incrementalmente sin reescribir toda la instantánea (snapshot) cada vez.
Dicho de otro modo, BaseCheckpointSaver es el contrato entre PregelLoop y el backend de almacenamiento concreto: el motor sólo conoce esta interfaz, y cambiar el backend no exige tocar el motor; quien escribe su propio saver sólo tiene que rellenar esos cinco métodos para que el motor herede automáticamente las características de concurrencia y persistencia de ese backend.
Motivación de diseño
- Estado y escrituras por separado:
putguarda la instantánea completa,put_writesguarda las escrituras incrementales de una tarea en un paso (put_writes:300-318). La separación permite que, al recuperarse de una interrupción, se reconstruya con precisión «qué tarea ya escribió en este paso y cuál no», en lugar de sólo poder retroceder al checkpoint anterior. thread_idcomo clave primaria: el config que devuelveputsólo llevathread_id/checkpoint_ns/checkpoint_id(put:277-298); entre llamadas y entre procesos, un checkpoint se localiza con estos tres campos.checkpoint_nses el namespace para subgrafos; la cadena vacía corresponde al grafo de nivel superior.- Interfaces síncrona y asíncrona: cada método de escritura tiene una versión async
aput/aget/alist/aput_writes(async methods:417-509), cuya implementación por defecto lanzaNotImplementedError. Una subclase de saver puede implementar sólo la versión síncrona y dejar que la versión async delegue conasyncio.to_thread, o escribir una ruta async nativa con un driver nativo. - Metadatos con etiqueta de origen: el campo
sourcedeCheckpointMetadatatoma los valoresinput/loop/update/fork(source:41-48).listy el filtrado se apoyan en los metadatos, así que el saver debe persistirlos junto con el checkpoint, sin perderlos. - Bypass de DeltaChannel: las docstrings de
prune/delete_for_runs/copy_threadadvierten reiteradamente (DeltaChannel-aware ops:320-415) que, si el grafo usaDeltaChannel, un saver personalizado no puede limitarse a borrar el checkpoint más reciente: debe conservar además la cadena decheckpoint_writesde los ancestros; de lo contrario la reconstrucción delta devolverá silenciosamente vacío.
Archivos clave
BaseCheckpointSaver 类:176— la clase abstracta en sí; define las interfaces put/get/list/put_writes.Checkpoint TypedDict:92-123— la forma de los campos de un checkpoint: channel_values/channel_versions/versions_seen/id/ts.CheckpointTuple:139-146— el valor de retorno deget_tuple; empaqueta checkpoint, metadata, parent_config y pending_writes en un NamedTuple.CheckpointMetadata:38-86—source/step/parents/run_id/counters_since_delta_snapshotpara DeltaChannel.get 方法:227-237—getdelega por defecto enget_tupley devuelve sólo el campocheckpoint.put 方法:277-298— persiste el checkpoint completo + metadata y devuelve el config actualizado.put_writes 方法:300-318— persiste la lista de escrituras incrementales de una tarea bajo un checkpoint.delete_thread / delete_for_runs:320-348— limpieza a nivel de hilo (thread), con advertencias sobre DeltaChannel.WRITES_IDX_MAP:795— los cuatro canales de controlERROR/SCHEDULED/INTERRUPT/RESUMEusan índices negativos para reservar la primary key y evitar colisiones con las escrituras de tareas reales.
Flujo de datos
Al terminar cada superpaso, PregelLoop invoca al saver para persistir los valores actuales de todos los canales y los nuevos números de versión de este paso. La firma por defecto de put y la implementación de InMemorySaver lo dejan claro: el motor llama a put, el saver descompone channel_values por (thread_id, ns, channel, version) para guardar blobs, el cuerpo principal va a la tabla principal y se devuelve un nuevo config que apunta al checkpoint_id recién persistido.
def put(
self,
config: RunnableConfig,
checkpoint: Checkpoint,
metadata: CheckpointMetadata,
new_versions: ChannelVersions,
) -> RunnableConfig:
"""Store a checkpoint with its configuration and metadata."""
raise NotImplementedErrorEsta es la interfaz pura en la clase base (put method:277-298), sin implementación por defecto. Cada subclase decide qué campos van a la tabla principal, cuáles a blobs, si usar ON CONFLICT DO UPDATE, etc. Al leer de vuelta, get_tuple se encarga de deserializar los blobs y recomponer channel_values, e incluye los pending_writes del mismo checkpoint dentro del CheckpointTuple (get_tuple:239-251).
Límites y fallos
put_writesusa índices negativos para los canales de control:WRITES_IDX_MAPmapeaERROR/SCHEDULED/INTERRUPT/RESUMEa -1 / -2 / -3 / -4 (WRITES_IDX_MAP:795); la primary key(task_id, idx)de un saver personalizado debe aceptar valores negativos, o las escrituras de control chocarán con la primary key.getdelega enget_tuple: basta con que la subclase implementeget_tuple;getpor defecto te devuelve.checkpoint(get:227-237). Pero siget_tupleno se implementa, lanzaNotImplementedError; no puede devolverNonesilenciosamente.- async por defecto es
NotImplementedError: cuando un saver síncrono no implementaaput, la ruta async por defecto no conmuta automáticamente a un hilo —la subclase debe escribir su propioasyncio.to_thread, o el llamador debe usar la ruta síncrona (async defaults:468-509). metadatano se puede perder: el parámetrofilterdelistfiltra por campos de metadatos (list method:253-275); si un saver no almacena o almacena mal los metadatos, el filtrado y la paginación se quedan mudos.- Bajo DeltaChannel no se puede borrar de forma naïve: la docstring de
pruneen la clase base lo deja claro (prune:374-415) —quedarse sólo con el checkpoint más reciente hace que la reconstrucción del canal delta devuelva vacío sin error. Un saver personalizado debe recorrer la cadena de ancestros antes de borrar. checkpoint_nspor defecto es la cadena vacía: las llamadas desde un subgrafo llevan un ns no vacío; las de nivel superior usan"". La primary key del saver debe incluir el ns, o checkpoints homónimos de subgrafos distintos se sobrescribirán entre sí.
Resumen
BaseCheckpointSaver fija de una sola vez «cómo almacenar, leer, listar y escribir incrementos de un checkpoint» en cinco métodos + un TypedDict; el motor Pregel sólo habla con esta interfaz, y el backend se puede sustituir a voluntad. Dos implementaciones de referencia se describen en las páginas hermanas InMemorySaver + SqliteSaver y PostgresSaver; el ritmo de persistencia en runtime se trata en motor Pregel y PregelLoop.
Véase la documentación oficial: documentación de LangGraph · README