Skip to content

BaseCheckpointSaver : la forme de l'interface de persistance

源码版本1.2.9

Responsabilités

Dans le runtime de LangGraph, à chaque superpas (superstep) terminé, le moteur Pregel doit stocker les valeurs de tous ses canaux (channel) ainsi que les écritures (writes) produites lors de ce superpas, afin qu'un prochain appel avec le même thread_id puisse reprendre depuis l'état précédent, mais aussi pour reprendre après une interruption (interrupt), rejouer l'historique, ou faire du débogage par voyage dans le temps. BaseCheckpointSaver est l'interface abstraite de cette couche de persistance (BaseCheckpointSaver:176). Elle ne se préoccupe pas du lieu de stockage — dict en mémoire, SQLite, Postgres, Redis conviennent tous — et fixe seulement les méthodes qu'un saver doit implémenter : put / get / get_tuple / list / put_writes (put / get / list / put_writes:227-318).

Elle fige aussi la forme du concept de point de contrôle (checkpoint) à travers le TypedDict Checkpoint : Checkpoint TypedDict:92-123, qui contient channel_values, channel_versions, versions_seen, id / ts / v. Ce que put reçoit n'est pas seulement un Checkpoint : il y a aussi metadata et new_versions, et le saver décide quels champs vont dans la table principale et lesquels dans la table blob. put_writes stocke séparément « les écritures produites par cette tâche dans ce superpas », découplées du checkpoint, de sorte que les écritures puissent être persistées incrémentalement sans réécrire toute la snapshot à chaque écriture.

Autrement dit, BaseCheckpointSaver est le contrat entre PregelLoop et le backend de stockage : le moteur ne connaît que cette interface, et un backend peut changer d'implémentation sans toucher au moteur ; réciproquement, il suffit qu'un backend implémente les cinq méthodes pour que le moteur hérite automatiquement de ses propriétés de concurrence et de persistance.

Motivation de conception

  • État et écritures stockés séparément : put stocke la snapshot complète, put_writes stocke les écritures incrémentales d'une tâche dans un superpas (put_writes:300-318). Cette séparation permet, lors d'une reprise après interruption, de reconstituer précisément « quelle tâche a déjà écrit dans ce superpas, laquelle non », plutôt que de devoir revenir globalement au checkpoint précédent.
  • thread_id comme clé primaire : le config renvoyé par put ne contient que thread_id / checkpoint_ns / checkpoint_id (put:277-298) ; ces trois champs localisent un checkpoint à travers les appels et les processus. checkpoint_ns est le namespace utilisé pour les sous-graphes ; la chaîne vide correspond au graphe de premier niveau.
  • Deux jeux d'interfaces, sync et async : chaque méthode d'écriture a son équivalent async aput / aget / alist / aput_writes (async methods:417-509), dont l'implémentation par défaut lève NotImplementedError. Une sous-classe peut n'implémenter que la version synchrone ; la version async se rabat par défaut sur asyncio.to_thread, ou peut fournir une vraie route async via un driver natif.
  • Métadonnées avec étiquette de source : dans CheckpointMetadata, le champ source prend les valeurs input / loop / update / fork (source:41-48). list et le filtrage reposent sur metadata : le saver doit donc stocker metadata avec le checkpoint, sans perte.
  • Bypass DeltaChannel : les docstrings de prune / delete_for_runs / copy_thread mettent en garde (DeltaChannel-aware ops:320-415) : si le graphe utilise DeltaChannel, un saver personnalisé ne peut pas se contenter de supprimer le dernier checkpoint ; il doit préserver la chaîne d'ancêtres des checkpoint_writes, sinon la reconstruction delta renverra silencieusement du vide.

Fichiers clés

  • BaseCheckpointSaver 类:176 — la classe abstraite elle-même, qui fixe les interfaces put/get/list/put_writes.
  • Checkpoint TypedDict:92-123 — la forme d'un checkpoint : channel_values/channel_versions/versions_seen/id/ts.
  • CheckpointTuple:139-146 — valeur de retour de get_tuple : un NamedTuple regroupant checkpoint, metadata, parent_config, pending_writes.
  • CheckpointMetadata:38-86source / step / parents / run_id / counters_since_delta_snapshot pour DeltaChannel.
  • get 方法:227-237get délègue par défaut à get_tuple et ne renvoie que le champ checkpoint.
  • put 方法:277-298 — stocke le checkpoint complet + metadata, renvoie le config mis à jour.
  • put_writes 方法:300-318 — stocke la liste des écritures incrémentales d'une tâche sous un checkpoint.
  • delete_thread / delete_for_runs:320-348 — nettoyage au niveau du thread, avec avertissement DeltaChannel.
  • WRITES_IDX_MAP:795ERROR / SCHEDULED / INTERRUPT / RESUME, les quatre canaux de contrôle, utilisent des idx négatifs comme sentinelle pour éviter de collisionner avec les écritures réelles sur la clé primaire.

Flux de données

À chaque superpas terminé, PregelLoop appelle le saver pour persister les valeurs courantes de tous les canaux et les nouveaux numéros de version de ce superpas. La signature par défaut de put et l'implémentation de InMemorySaver suffisent à rendre cela clair : le moteur appelle put, le saver éclate les channel_values en blobs stockés par (thread_id, ns, channel, version), met le reste dans la table principale, et renvoie un nouveau config pointant vers le checkpoint_id qui vient d'être écrit.

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

C'est une interface pure dans la classe de base (put method:277-298) : pas d'implémentation par défaut. Chaque sous-classe décide quels champs vont dans la table principale, lesquels dans blob, et si elle utilise ON CONFLICT DO UPDATE ou similaire. À la lecture, get_tuple est chargé de désérialiser les blobs pour reconstruire channel_values, et d'ajouter les pending_writes du même checkpoint dans CheckpointTuple (get_tuple:239-251).

Limites et échecs

  • put_writes utilise des idx négatifs pour les canaux de contrôle : WRITES_IDX_MAP mappe ERROR / SCHEDULED / INTERRUPT / RESUME vers -1 / -2 / -3 / -4 (WRITES_IDX_MAP:795). La clé primaire (task_id, idx) d'un saver personnalisé doit accepter les valeurs négatives, sinon les écritures de contrôle collisionnent.
  • get appelle par défaut get_tuple : une sous-classe n'a besoin d'implémenter que get_tuple ; get extrait .checkpoint pour vous (get:227-237). Mais si get_tuple n'est pas implémentée, elle lève NotImplementedError — pas de retour silencieux à None.
  • async par défaut à NotImplementedError : si un saver sync n'implémente pas aput, la route async ne bascule pas automatiquement en thread — soit la sous-classe écrit elle-même le pont asyncio.to_thread, soit l'appelant passe par la route synchrone (async defaults:468-509).
  • metadata ne doit pas être perdu : le paramètre filter de list filtre justement sur les champs de metadata (list method:253-275) ; si le saver ne stocke pas metadata ou la stocke mal, le filtrage et la pagination deviennent muets.
  • Pas de suppression naïve sous DeltaChannel : le docstring de prune le dit explicitement (prune:374-415) — ne garder que le dernier checkpoint ferait que la reconstruction delta renvoie silencieusement du vide, sans erreur. Un saver personnalisé doit remonter la chaîne d'ancêtres avant de supprimer.
  • checkpoint_ns par défaut est chaîne vide : un appel de sous-graphe porte un ns non vide, un appel de premier niveau a "". La clé primaire du saver doit intégrer ns, sinon des sous-graphes distincts ayant le même nom de checkpoint s'écrasent mutuellement.

Résumé

BaseCheckpointSaver fige en une fois « comment un checkpoint est stocké, lu, listé et comment ses écritures incrémentales sont persistées » : cinq méthodes + un TypedDict. Le moteur Pregel ne parle qu'à cette interface, et le backend est libre d'évoluer. Deux implémentations de référence sont détaillées dans les pages sœurs InMemorySaver + SqliteSaver et PostgresSaver ; le rythme de persistance à l'exécution est expliqué dans moteur Pregel et PregelLoop.

Voir la documentation officielle : documentation LangGraph · README