Skip to content

BaseCheckpointSaver: la forma de la interfaz de persistencia

源码版本1.2.9

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: put guarda la instantánea completa, put_writes guarda 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_id como clave primaria: el config que devuelve put sólo lleva thread_id / checkpoint_ns / checkpoint_id (put:277-298); entre llamadas y entre procesos, un checkpoint se localiza con estos tres campos. checkpoint_ns es 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 lanza NotImplementedError. Una subclase de saver puede implementar sólo la versión síncrona y dejar que la versión async delegue con asyncio.to_thread, o escribir una ruta async nativa con un driver nativo.
  • Metadatos con etiqueta de origen: el campo source de CheckpointMetadata toma los valores input / loop / update / fork (source:41-48). list y 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_thread advierten reiteradamente (DeltaChannel-aware ops:320-415) que, si el grafo usa DeltaChannel, un saver personalizado no puede limitarse a borrar el checkpoint más reciente: debe conservar además la cadena de checkpoint_writes de 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 de get_tuple; empaqueta checkpoint, metadata, parent_config y pending_writes en un NamedTuple.
  • CheckpointMetadata:38-86source / step / parents / run_id / counters_since_delta_snapshot para DeltaChannel.
  • get 方法:227-237get delega por defecto en get_tuple y devuelve sólo el campo checkpoint.
  • 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 control ERROR / SCHEDULED / INTERRUPT / RESUME usan í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.

python
def put(
    self,
    config: RunnableConfig,
    checkpoint: Checkpoint,
    metadata: CheckpointMetadata,
    new_versions: ChannelVersions,
) -> RunnableConfig:
    """Store a checkpoint with its configuration and metadata."""
    raise NotImplementedError

Esta 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_writes usa índices negativos para los canales de control: WRITES_IDX_MAP mapea ERROR / SCHEDULED / INTERRUPT / RESUME a -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.
  • get delega en get_tuple: basta con que la subclase implemente get_tuple; get por defecto te devuelve .checkpoint (get:227-237). Pero si get_tuple no se implementa, lanza NotImplementedError; no puede devolver None silenciosamente.
  • async por defecto es NotImplementedError: cuando un saver síncrono no implementa aput, la ruta async por defecto no conmuta automáticamente a un hilo —la subclase debe escribir su propio asyncio.to_thread, o el llamador debe usar la ruta síncrona (async defaults:468-509).
  • metadata no se puede perder: el parámetro filter de list filtra 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 prune en 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_ns por 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