Skip to content

InMemorySaver + SqliteSaver : checkpoints in-process et au niveau fichier

源码版本1.2.9

Responsabilités

BaseCheckpointSaver fixe la forme de l'interface ; les deux savers les plus courants qui l'incarnent sont InMemorySaver et SqliteSaver — le premier stuffe tout l'état dans un defaultdict en mémoire du processus, le second éclate les mêmes champs en deux tables SQLite checkpoints et writes. Ensemble, ils couvrent le débogage et la persistance mono-machine ; en production on passe à PostgresSaver.

InMemorySaver vit dans libs/checkpoint/langgraph/checkpoint/memory/__init__.py (InMemorySaver:33) ; en interne, trois defaultdict stockent trois choses : storage pour le corps du checkpoint + métadonnées + parent_id, writes pour (thread_id, ns, checkpoint_id) → (task_id, idx) → enregistrement d'écriture, et blobs pour (thread_id, ns, channel, version) → binaire sérialisé (storage / writes / blobs:68-83). Cette structure à trois compartiments correspond un à un à la structure trois-tables de Postgres/SQLite, sauf qu'un dict remplace les tables.

SqliteSaver vit dans libs/checkpoint-sqlite/langgraph/checkpoint/sqlite/__init__.py (SqliteSaver:45). Il hérite directement de BaseCheckpointSaver[str], prend un sqlite3.Connection, et dans setup() crée deux tables (<SrcLink path="libs/checkpoint-sqlite/langgraph/checkpoint/sqlite/__init__.py" lines="129-166" label="setup"/>). put / put_writes utilisent le SQL standard INSERT OR REPLACE / INSERT OR IGNORE. Il embarque un threading.Lock pour sérialiser la connexion unique, donc check_same_thread=False est sûr.

À noter : MemorySaver = InMemorySaver est un alias historique (MemorySaver alias:631) ; l'import classique from langgraph.checkpoint.memory import MemorySaver donne exactement la même classe.

Motivation de conception

  • Les trois dicts d'InMemory reflètent les trois tables : storage / writes / blobs (data members:68-83) correspondent aux tables checkpoints / writes de SQLite et aux trois tables checkpoints / checkpoint_writes / checkpoint_blobs de Postgres. InMemory isole aussi les blobs dans un dict à part pour que la sérialisation d'une valeur de canal n'arrive qu'une fois par version, plutôt que de tout resérialiser à chaque put du checkpoint complet.
  • PersistentDict pour une persistance optionnelle : le module memory expose aussi PersistentDict(defaultdict) (PersistentDict:634), qui dump le dict en pickle vers un fichier et le recharge au redémarrage. InMemorySaver.__init__ accepte un paramètre factory (__init__:85-99) précisément pour que factory=PersistentDict branche automatiquement le context manager — un compromis pour « débogage, mais je veux survivre aux redémarrages ».
  • SQLite en WAL + un seul lock : la première instruction de setup() est PRAGMA journal_mode=WAL (WAL:141) pour que les lectures ne bloquent pas les écritures ; mais les écritures restent sérialisées par self.lock car sqlite3.Connection n'est pas thread-safe. Le docstring prévient explicitement « ne scale pas en multi-thread » (note:48-54), pour la concurrence il faut passer à Postgres.
  • La table writes distingue REPLACE et IGNORE : put_writes en SQLite choisit entre INSERT OR REPLACE et INSERT OR IGNORE selon que toutes les écritures tombent dans WRITES_IDX_MAP ou non (put_writes query:462-466) : les écritures sur canaux de contrôle (ERROR / INTERRUPT / RESUME) peuvent écraser, les écritures métier ordinaires sont ignorées en cas de collision — sémantique identique à if inner_key[1] >= 0 ... continue dans InMemory (InMemorySaver.put_writes:499-509).
  • max(checkpoints.keys()) comme « dernier » pour InMemory : pas de ORDER BY explicite, on s'appuie sur la croissance monotone de la chaîne checkpoint_id (latest:282-283) ; checkpoint_id doit donc être un UUID6/UUID7 comparable lexicographiquement, pas un UUID aléatoire pur.

Fichiers clés

Flux de données

Le bloc InMemorySaver.put ci-dessous réalise en une fois « éclater channel_values dans blobs » et « mettre le corps dans storage », et montre très clairement comment le saver décompose un Checkpoint TypedDict en structure trois niveaux :

python
c = checkpoint.copy()
thread_id = config["configurable"]["thread_id"]
checkpoint_ns = config["configurable"]["checkpoint_ns"]
values: dict[str, Any] = c.pop("channel_values")  # type: ignore[misc]
for k, v in new_versions.items():
    self.blobs[(thread_id, checkpoint_ns, k, v)] = (
        self.serde.dumps_typed(values[k]) if k in values else ("empty", b"")
    )
self.storage[thread_id][checkpoint_ns].update(
    {
        checkpoint["id"]: (
            self.serde.dumps_typed(c),
            self.serde.dumps_typed(get_checkpoint_metadata(config, metadata)),
            config["configurable"].get("checkpoint_id"),  # parent
        )
    }
)

(put body:448-464) Notez que blobs est indexé par (thread_id, checkpoint_ns, channel, version), tandis que storage est indexé par (thread_id, checkpoint_ns, checkpoint_id) — deux indexations disjointes : la valeur du canal est découplée du corps du checkpoint, et un canal dont la version n'a pas changé n'est pas resérialisé. À la lecture, get_tuple appelle _load_blobs pour désérialiser, pour chaque (channel, version) de channel_versions, le blob correspondant et reconstruire channel_values (_load_blobs:259-265) puis l'empaqueter dans CheckpointTuple.

SqliteSaver suit la même logique, mais le dict devient SQL, et le compartiment blobs devient la colonne checkpoint BLOB de la table principale + la colonne value BLOB de la table writes — il n'y a pas de table blob dédiée, toutes les valeurs de canaux sont directement sérialisées dans le champ checkpoint blob :

python
cur.execute(
    "INSERT OR REPLACE INTO checkpoints (thread_id, checkpoint_ns, checkpoint_id, parent_checkpoint_id, type, checkpoint, metadata) VALUES (?, ?, ?, ?, ?, ?, ?)",
    (
        str(config["configurable"]["thread_id"]),
        checkpoint_ns,
        checkpoint["id"],
        config["configurable"].get("checkpoint_id"),
        type_,
        serialized_checkpoint,
        serialized_metadata,
    ),
)

(SqliteSaver.put INSERT:425-436)

Limites et échecs

  • InMemorySaver perd tout si le processus meurt : le docstring le dit clairement — « n'utilisez InMemorySaver que pour le débogage ou les tests » (warning:40-44). Pour persister, passer à Sqlite ou Postgres.
  • InMemorySaver trouve le dernier via l'ordre lexicographique des checkpoint_id : max(checkpoints.keys()) (max:282-283) suppose un ID monotone ; si un saver personnalisé utilise un ID non monotone, get_tuple sans checkpoint_id renverra la mauvaise ligne.
  • SqliteSaver : connexion unique, lock unique : le context manager cursor() prend self.lock (cursor:181-189), les écritures concurrentes sont sérialisées et une longue transaction bloque les autres lectures. Pour la concurrence, il faut le mode pipeline de Postgres.
  • SqliteSaver utilise des ? avec un ordre de colonnes codé en dur : l'ordre des colonnes doit correspondre à la liste INSERT INTO checkpoints (...) (INSERT columns:426) ; ajouter une colonne en migration implique de modifier les deux côtés ; la colonne task_path a été ajoutée a posteriori, voir la fin de MIGRATIONS dans Postgres (task_path migration:90).
  • AsyncSqliteSaver n'est pas qu'un wrapper to_thread : AsyncSqliteSaver:38 est une classe distincte, dont aput / aget_tuple / alist sont réécrits avec aiosqlite (aput:509) — ce n'est pas le pont asyncio.to_thread par défaut de la classe de base, donc la route async gère vraiment la concurrence sans être sérialisée par le lock global.
  • Avec factory=PersistentDict, il faut appeler close() en sortie : PersistentDict.sync() dump le pickle vers le fichier (sync:657) ; sans appel à __exit__, rien n'est persisté ; with InMemorySaver(factory=PersistentDict, ...) est l'idiome recommandé.

Résumé

InMemorySaver est l'implémentation de référence « trois dicts », SqliteSaver déploie les trois compartiments en deux tables SQLite + WAL ; ensemble ils couvrent le débogage et la persistance mono-machine. Les conventions de clé primaire, la séparation des blobs, les idx négatifs pour les canaux de contrôle trouvent ici leur forme concrète. Pour la production à l'échelle, voir PostgresSaver ; pour la forme de l'interface, voir BaseCheckpointSaver.

Voir la documentation officielle : documentation LangGraph · README