BaseCheckpointSaver: Die Form der Persistenz-Schnittstelle
Verantwortung
In der LangGraph-Laufzeit muss die Pregel-Engine nach jedem Superstep (superstep) alle ihre aktuellen Kanalwerte (channel) zusammen mit den in diesem Schritt erzeugten Writes (writes) speichern, damit beim nächsten Eintreffen desselben thread_id die Berechnung am vorherigen Zustand fortgesetzt werden kann, und auch, um nach einer Unterbrechung (interrupt) den Zustand wiederherzustellen, den Verlauf zu replayen und Time-Travel-Debugging zu betreiben. BaseCheckpointSaver ist genau diese Abstraktionsschicht für Persistenz (BaseCheckpointSaver:176). Es kümmert sich nicht darum, wohin die Daten gespeichert werden – In-Memory-dict, SQLite, Postgres, Redis alles möglich – sondern legt nur die Methoden fest, die ein Saver implementieren muss: put / get / get_tuple / list / put_writes (put / get / list / put_writes:227-318).
Gleichzeitig fixiert es die Form des Konzepts Checkpoint (checkpoint) als Checkpoint-TypedDict: Checkpoint TypedDict:92-123, das channel_values, channel_versions, versions_seen, id / ts / v enthält. put empfängt nicht nur einen Checkpoint, sondern auch metadata und new_versions; der Saver entscheidet, welche Felder in die Haupttabelle und welche in die Blob-Tabelle wandern. put_writes spe separat die Schreibvorgänge, die ein Task in einem Schritt erzeugt – entkoppelt vom eigentlichen Checkpoint, sodass Writes inkrementell persistiert werden können, ohne bei jedem Write den gesamten Snapshot neu zu schreiben.
Anders ausgedrückt: BaseCheckpointSaver ist der Vertrag zwischen PregelLoop und dem konkreten Speicher-Backend; die Engine spricht nur diese Schnittstelle, das Backend kann ausgetauscht werden, ohne die Engine zu ändern; schreibt ein Backend seinen eigenen Saver, indem es die fünf Methoden ausfüllt, erbt die Engine automatisch dessen Nebenläufigkeits- und Persistenz-Eigenschaften.
Entwurfsmotivation
- Zustand und Writes getrennt speichern:
putspeichert den gesamten Snapshot,put_writesspeichert die inkrementellen Schreibvorgänge eines einzelnen Tasks in einem Schritt (put_writes:300-318). Die Trennung erlaubt es, bei der Wiederherstellung nach Unterbrechung exakt «welcher Task in diesem Schritt bereits geschrieben hat und welcher nicht» zu rekonstruieren, statt nur global auf den vorherigen Checkpoint zurückzurollen. thread_idals Primärschlüssel: Das vonputzurückgegebene Config enthält nurthread_id/checkpoint_ns/checkpoint_id(put:277-298); über Aufrufe und Prozesse hinweg wird ein Checkpoint über diese drei Felder lokalisiert.checkpoint_nsist der Namensraum für Untergraphen (subgraph), eine leere Zeichenkette bezeichnet den Top-Level-Graphen.- Synchrone + asynchrone Doppel-Schnittstelle: Jede Schreibmethode hat eine async-Version
aput/aget/alist/aput_writes(async methods:417-509), deren StandardimplementierungNotImplementedErrorwirft. Eine Saver-Subklasse kann nur die synchrone Variante implementieren; die asynchrone fällt dann aufasyncio.to_threadzurück, oder sie schreibt einen echten async-Pfad mit einem nativen Treiber. - Metadata mit Quellen-Tag: Das Feld
sourceinCheckpointMetadatanimmtinput/loop/update/fork(source:41-48). list und Filterung arbeiten beide über metadata, daher muss der Saver metadata zusammen mit dem Checkpoint speichern und darf es nicht verwerfen. - DeltaChannel-Seitenweg: Die Docstrings von
prune/delete_for_runs/copy_threadwarnen wiederholt (DeltaChannel-aware ops:320-415): Wenn der GraphDeltaChannelverwendet, darf ein eigener Saver nicht einfach den neuesten Checkpoint löschen – diecheckpoint_writesauf der Ahnenkette müssen erhalten bleiben, sonst liefert der Delta-Wiederaufbau stillschweigend leere Ergebnisse.
Schlüsseldateien
BaseCheckpointSaver Klasse:176— Die abstrakte Basisklasse, legt die put/get/list/put_writes-Schnittstellen fest.Checkpoint TypedDict:92-123— Felder eines Checkpoints: channel_values/channel_versions/versions_seen/id/ts.CheckpointTuple:139-146— Rückgabewert vonget_tuple, fasst checkpoint, metadata, parent_config und pending_writes in einem NamedTuple zusammen.CheckpointMetadata:38-86—source/step/parents/run_id/ für DeltaChannelcounters_since_delta_snapshot.get Methode:227-237—getdelegiert standardmäßig anget_tupleund gibt nur dascheckpoint-Feld zurück.put Methode:277-298— Speichert einen vollständigen Checkpoint inkl. metadata und gibt das aktualisierte Config zurück.put_writes Methode:300-318— Speichert die Liste inkrementeller Schreibvorgänge eines Tasks zu einem gegebenen Checkpoint.delete_thread / delete_for_runs:320-348— Thread-Level Aufräumung mit DeltaChannel-Warnung.WRITES_IDX_MAP:795—ERROR/SCHEDULED/INTERRUPT/RESUMEbelegen als Kontroll-Channels negative Indizes, um Kollisionen mit echten Task-Writes beim Primärschlüssel zu vermeiden.
Datenfluss
Nach jedem Superstep ruft PregelLoop den Saver auf, um die aktuellen Kanalwerte und die neuen Versionsnummern dieses Schritts zu persistieren. Die Signatur von put und die Implementierung von InMemorySaver reichen aus, um das klar zu machen – die Engine ruft put auf, der Saver zerlegt channel_values nach (thread_id, ns, channel, version) und speichert Blob und Hauptkörper in der Haupttabelle; das neue Config zeigt auf die gerade geschriebene checkpoint_id.
def put(
self,
config: RunnableConfig,
checkpoint: Checkpoint,
metadata: CheckpointMetadata,
new_versions: ChannelVersions,
) -> RunnableConfig:
"""Store a checkpoint with its configuration and metadata."""
raise NotImplementedErrorDas ist die reine Schnittstelle in der Basisklasse (put method:277-298), ohne Standardimplementierung. Jede Subklasse entscheidet selbst, welche Felder in die Haupttabelle bzw. in Blobs wandern und ob ON CONFLICT DO UPDATE oder ähnliches verwendet wird. Beim Zurücklesen muss get_tuple die Blobs deserialisieren und zurück in channel_values einsetzen sowie die pending_writes desselben Checkpoints in das CheckpointTuple packen (get_tuple:239-251).
Grenzen und Fehler
put_writesnutzt negative Indizes für Kontroll-Channels:WRITES_IDX_MAPbildetERROR/SCHEDULED/INTERRUPT/RESUMEauf -1 / -2 / -3 / -4 ab (WRITES_IDX_MAP:795); der Primärschlüssel(task_id, idx)eines eigenen Savers muss negative Werte akzeptieren, sonst kollidieren Kontroll-Writes.getruft standardmäßigget_tupleauf: Eine Subklasse braucht nurget_tuplezu implementieren,getliefert per Default.checkpoint(get:227-237). Werget_tuplenicht überschreibt, erhältraise NotImplementedError, nicht stillschweigendNone.- async fällt standardmäßig auf
NotImplementedError: Wenn ein synchroner Saveraputnicht implementiert, leitet der Default-Async-Pfad nicht automatisch in einen Thread um – entweder schreibt die Subklasse selbstasyncio.to_thread, oder der Aufrufer geht über den synchronen Pfad (async defaults:468-509). metadatadarf nicht verloren gehen: Derfilter-Parameter vonlistfiltert nach metadata-Feldern (list method:253-275); ein Saver, der metadata nicht oder fehlerhaft speichert, bringt Filterung und Paginierung zum Verstummen.- Unter DeltaChannel nicht naiv löschen: Der Docstring von
prunein der Basisklasse sagt ausdrücklich (prune:374-415) – nur den neuesten Checkpoint zu behalten, führt dazu, dass der Delta-Channel-Wiederaufbau stillschweigend leer zurückkehrt, ohne Fehler. Ein eigener Saver muss vor dem Löschen die Ahnenkette durchlaufen. checkpoint_nsist standardmäßig leer: Subgraph-Aufrufe tragen einen nicht-leeren ns, Top-Level-Aufrufe verwenden"". Der Primärschlüssel des Savers muss ns einbeziehen, sonst überschreiben sich gleichnamige Checkpoints verschiedener Untergraphen.
Zusammenfassung
BaseCheckpointSaver fixiert «wie ein Checkpoint gespeichert, gelesen, gelistet und inkrementell geschrieben wird» ein für alle Mal als fünf Methoden plus ein TypedDict; die Pregel-Engine spricht nur mit der Schnittstelle, das Backend kann beliebig ausgetauscht werden. Zwei konkrete Referenzimplementierungen finden sich auf den Schwesterseiten InMemorySaver + SqliteSaver und PostgresSaver; wie die Laufzeit den Takt der Persistenz bestimmt, siehe Pregel 引擎 und PregelLoop.
Siehe offizielle Dokumentation: LangGraph 文档 · README.