StateSnapshot: Zustandssnapshot nach jedem Superstep
Verantwortung
StateSnapshot ist die unveränderliche Momentaufnahme, die LangGraph an jeder Superstep (superstep)-Grenze nach außen exposes: «Wie sieht der gesamte Graph aus, nachdem dieser Schritt gerade gelaufen ist». Es packt fünf Dinge in ein einziges NamedTuple: die aktuellen Werte jedes Kanals (channel) (values), welche Knoten als Nächstes laufen sollen (next), die in diesem Schritt verwendete config, die Metadaten des zugehörigen Checkpoints (metadata + created_at) und die config des Eltern-Checkpoints(StateSnapshot 定义:643-661). Was graph.get_state(config), graph.aget_state(config) oder eine Iteration über graph.get_state_history(config) zurückgeben, ist genau dieser Typ.
In der Architektur liegt es «über dem Checkpoint (checkpoint), unter der Benutzeransicht». Der Checkpoint (CheckpointTuple) ist das interne Serialisierungsformat des Engines und speichert Low-Level-Felder wie channel_values / channel_versions / pending_writes direkt — man sieht nicht, wie der Geschäftsstatus aussieht; StateSnapshot übersetzt diese Schicht in «ein dict, das Geschäfte direkt lesen können + was als Nächstes läuft + Aufgabenliste»(_prepare_state_snapshot:1145-1162). Auch der Zustand von Untergraphen geht darüber — das Feld task.state enthält bei rekursivem subgraphs=True den StateSnapshot eines Untergraphen(PregelTask.state:200-201).
StateSnapshot selbst ist ein NamedTuple und unveränderlich; um den «Zustand zu ändern», muss man über graph.update_state(config, values) neue Werte schreiben. Das Engine wandelt das in einen pending write mit NULL_TASK_ID um, der im nächsten Checkpoint landet, und erzeugt anschließend ein neues StateSnapshot.
Entwurfsmotivation
Warum CheckpointTuple nicht direkt den Benutzern expose?
- Interne Darstellung entkoppeln:Das Checkpoint-Format speichert zur effizienten Serialisierung ein flaches dict mit Versionsnummern wie
channel_versions; die Benutzer interessieren sich aber für «mein Geschäfts-state-dict», weshalbread_channelsdie Kanalwerte auslesen muss(read_channels:1258). Eine zusätzlicheStateSnapshot-Schicht erlaubt dem Checkpoint-Format, sich frei weiterzuentwickeln, ohne die externe API zu brechen. - «Nächster Schritt»-Semantik einbetten:Das Engine weiß, welche Knoten als Nächstes laufen, aber
CheckpointTupleselbst speichert nur die geschriebenen Werte;StateSnapshotruft beim Erzeugenprepare_next_tasks(prepare_next_tasks:1178-1195) auf und berechnetnext=(node_name, ...), was die Frage «Was passiert als Nächstes» direkt beantwortet. - Aufgabenliste sichtbar:
StateSnapshot.tasksist das Tupel derPregelTask, die im aktuellen Schritt ausgeführt werden(tasks:658), und enthältwrites/interrupts/ Zustand von Untergraphen — entscheidend für Human-in-the-Loop, um zu beurteilen «an welchem Knoten hängt es gerade und auf welche Eingabe wird gewartet». - Unveränderlich = vergleichbar:Historische Snapshots bilden über
parent_configeine verkettete Liste(parent_config:656);get_state_historyliefert die in umgekehrter zeitlicher Reihenfolge sortierte Snapshot-Liste. Jeder einzelne Snapshot kann über seineconfigfür Time-Travel neu referenziert werden. - Pending writes transparent anwenden:
get_statewendet standardmäßig die noch nicht in den Checkpoint geschriebenenpending_writestemporär auf die Kanäle an und liest sie dann(apply pending writes:1239-1249). Wenn Sieget_stateaufrufen, kurz nachdem der Knoten geschrieben hat, aber bevor der Checkpoint gesetzt wurde, sehen Sie den neuesten Wert und nicht den veralteten vom vorherigen Schritt.
Schlüsseldateien
StateSnapshot NamedTuple:643-661— 7 Felder:values/next/config/metadata/created_at/parent_config/tasks/interrupts.PregelTask:200-217— Aufgaben-Objekt; das Feldstatespeichert bei Rekursion denStateSnapshotdes Untergraphen oder eineRunnableConfig, die ihn lokalisiert._prepare_state_snapshot:1145-1162— Fallback bei leerem saved: liefert einen leeren Snapshot mitvalues={}/next=()/tasks=().step 计算 + prepare_next_tasks:1167-1195— leitet die aktuelle Schritt-nummer aussaved.metadata["step"]ab und ruftprepare_next_tasksauf, um das Aufgaben-dict dieses Schritts zu berechnen.task_states 递归:1199-1227— baut für jede Aufgabe mittask.name in subgraphseincheckpoint_nsund ruft rekursivget_statedes Untergraphen auf.apply pending writes:1229-1255— wendet Schreiboperationen mitNULL_TASK_IDund reguläre pending writes temporär auf die Kanäle an und berechnet danntasks_with_writes.组装 StateSnapshot:1257-1266— setzt Kanalwerte, next-Knoten-tuple, config, metadata, ts, parent_config, tasks, interrupts zum endgültigenStateSnapshotzusammen.get_state:1392-1434— synchroner Einstieg: löst den checkpointer auf, behandelt das Routing übercheckpoint_nszum Untergraphen und entscheidet überapply_pending_writes.get_state_history:1480— iteriert über alle historischen Checkpoints derselbenthread_idund erzeugt für jeden einStateSnapshot._aprepare_state_snapshot:1268— asynchrone Version; Logik symmetrisch zur synchronen, nurchannels_from_checkpointwird durchachannels_from_checkpointersetzt.
Datenfluss
Der folgende Abschnitt ist der Schlüssel, in dem _prepare_state_snapshot den Checkpoint in ein StateSnapshot übersetzt: Schritt-nummer aus saved.metadata ableiten, prepare_next_tasks aufrufen, um die Aufgaben dieses Schritts zu holen, und dann die Untergraphen-Rekursion behandeln.
step = saved.metadata.get("step", -1) + 1
stop = step + 2
channels, managed = channels_from_checkpoint(
self.channels,
saved.checkpoint,
saver=self.checkpointer
if isinstance(self.checkpointer, BaseCheckpointSaver) else None,
config=saved.config,
)
next_tasks = prepare_next_tasks(
saved.checkpoint,
saved.pending_writes or [],
self.nodes,
channels,
managed,
saved.config,
step,
stop,
for_execution=True,
store=self.store,
checkpointer=(self.checkpointer
if isinstance(self.checkpointer, BaseCheckpointSaver) else None),
manager=None,
)(prepare_next_tasks 调用:1167-1195)
Abschließend werden die ausgelesenen Kanalwerte und die gefilterten Aufgabennamen (ohne die mit bereits geschriebenen pending) zum unveränderlichen Snapshot zusammengesetzt:
return StateSnapshot(
read_channels(channels, self.stream_channels_asis),
tuple(t.name for t in next_tasks.values() if not t.writes),
patch_checkpoint_map(saved.config, saved.metadata),
saved.metadata,
saved.checkpoint["ts"],
patch_checkpoint_map(saved.parent_config, saved.metadata),
tasks_with_writes,
tuple([i for task in tasks_with_writes for i in task.interrupts]),
)Grenzen und Fehler
- Ohne Checkpoint wird ein leerer Snapshot zurückgegeben statt zu werfen(
空 saved:1152-1162),values={}/next=(). Die aufrufende Seite muss selbst prüfen, obsnapshot.nextein leeres Tupel ist, um zwischen «Graph noch nicht gelaufen» und «Graph fertig» zu unterscheiden. get_stateohne checkpointer wirft direktValueError(no checkpointer:1401-1402),No checkpointer set— ein rein im Speicher laufender Graph hat keine historischen Snapshots.- Wenn
checkpoint_nskeinen passenden Untergraphen findet, wird ein Fehler geworfen(subgraph not found:1415-1416),Subgraph {recast} not found; tritt häufig auf, wenn der Namespace fehlerhaft mitNS_SEP(|) oderNS_END(:) zusammengesetzt wurde. apply_pending_writeswird nur aktiviert, wenncheckpoint_idnicht explizit angegeben wurde(apply 条件:1433). D. h. ein historischer Snapshot zeigt genau die damals gespeicherten Werte und wird nicht durch spätere pending writes verunreinigt; nur der «neueste» Snapshot wendet pending temporär an.- Das Feld
nextfiltert Aufgaben heraus, die bereits writes haben(next 过滤:1259),Aufgaben mit nicht-leeremt.writeszählen nicht als «next» — da ihre Eingaben in diesem Schritt bereits konsumiert wurden, spiegeltnextdie «noch nicht gelaufenen Aufgaben» wider. tasksenthältinterrupts(interrupts 聚合:1265),dieInterrupt-Objekte aller Aufgaben werden zu einem flachen Tupel zusammengefasst; der Benutzer füllt sie in dieser Reihenfolge zurück, wenn erCommand(resume=...)aufruft.
Zusammenfassung
StateSnapshot ist der Grenztyp, mit dem LangGraph das interne Checkpoint-Format in eine geschäftslesbare Ansicht übersetzt — es speichert keine Daten, sondern setzt nur Checkpoint + pending writes plus «nächster Schritt»-Semantik zu einem unveränderlichen Snapshot zusammen. Alle APIs für Fortsetzen nach Haltepunkt, Time-Travel und Human-in-the-Loop bauen darauf auf.
Es bildet mit RunnableConfig ein Paar: config ist die Koordinate «wie dieser Snapshot zu finden ist», und StateSnapshot ist «was man sieht, nachdem man ihn gefunden hat». Weiter geht es mit StateGraph, wo die Form des Zustands definiert wird, der gesnapshotted wird, oder mit Pregel 引擎, wo beschrieben wird, wie _prepare_state_snapshot von get_state aufgerufen wird.
Siehe offizielle Dokumentation: LangGraph 文档 · README.