BaseCheckpointSaver : la forme de l'interface de persistance
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 :
putstocke la snapshot complète,put_writesstocke 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_idcomme clé primaire : le config renvoyé parputne contient quethread_id/checkpoint_ns/checkpoint_id(put:277-298) ; ces trois champs localisent un checkpoint à travers les appels et les processus.checkpoint_nsest 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èveNotImplementedError. Une sous-classe peut n'implémenter que la version synchrone ; la version async se rabat par défaut surasyncio.to_thread, ou peut fournir une vraie route async via un driver natif. - Métadonnées avec étiquette de source : dans
CheckpointMetadata, le champsourceprend les valeursinput/loop/update/fork(source:41-48).listet 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_threadmettent en garde (DeltaChannel-aware ops:320-415) : si le graphe utiliseDeltaChannel, un saver personnalisé ne peut pas se contenter de supprimer le dernier checkpoint ; il doit préserver la chaîne d'ancêtres descheckpoint_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 deget_tuple: un NamedTuple regroupant checkpoint, metadata, parent_config, pending_writes.CheckpointMetadata:38-86—source/step/parents/run_id/counters_since_delta_snapshotpour DeltaChannel.get 方法:227-237—getdélègue par défaut àget_tupleet ne renvoie que le champcheckpoint.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:795—ERROR/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.
def put(
self,
config: RunnableConfig,
checkpoint: Checkpoint,
metadata: CheckpointMetadata,
new_versions: ChannelVersions,
) -> RunnableConfig:
"""Store a checkpoint with its configuration and metadata."""
raise NotImplementedErrorC'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_writesutilise des idx négatifs pour les canaux de contrôle :WRITES_IDX_MAPmappeERROR/SCHEDULED/INTERRUPT/RESUMEvers -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.getappelle par défautget_tuple: une sous-classe n'a besoin d'implémenter queget_tuple;getextrait.checkpointpour vous (get:227-237). Mais siget_tuplen'est pas implémentée, elle lèveNotImplementedError— pas de retour silencieux àNone.- async par défaut à
NotImplementedError: si un saver sync n'implémente pasaput, la route async ne bascule pas automatiquement en thread — soit la sous-classe écrit elle-même le pontasyncio.to_thread, soit l'appelant passe par la route synchrone (async defaults:468-509). metadatane doit pas être perdu : le paramètrefilterdelistfiltre 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
prunele 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_nspar 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