Skip to content

StateSnapshot: Zustandssnapshot nach jedem Superstep

源码版本1.2.9

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», weshalb read_channels die Kanalwerte auslesen muss(read_channels:1258). Eine zusätzliche StateSnapshot-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 CheckpointTuple selbst speichert nur die geschriebenen Werte; StateSnapshot ruft beim Erzeugen prepare_next_tasks(prepare_next_tasks:1178-1195) auf und berechnet next=(node_name, ...), was die Frage «Was passiert als Nächstes» direkt beantwortet.
  • Aufgabenliste sichtbar:StateSnapshot.tasks ist das Tupel der PregelTask, die im aktuellen Schritt ausgeführt werden(tasks:658), und enthält writes / 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_config eine verkettete Liste(parent_config:656); get_state_history liefert die in umgekehrter zeitlicher Reihenfolge sortierte Snapshot-Liste. Jeder einzelne Snapshot kann über seine config für Time-Travel neu referenziert werden.
  • Pending writes transparent anwenden:get_state wendet standardmäßig die noch nicht in den Checkpoint geschriebenen pending_writes temporär auf die Kanäle an und liest sie dann(apply pending writes:1239-1249). Wenn Sie get_state aufrufen, 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 Feld state speichert bei Rekursion den StateSnapshot des Untergraphen oder eine RunnableConfig, die ihn lokalisiert.
  • _prepare_state_snapshot:1145-1162 — Fallback bei leerem saved: liefert einen leeren Snapshot mit values={} / next=() / tasks=().
  • step 计算 + prepare_next_tasks:1167-1195 — leitet die aktuelle Schritt-nummer aus saved.metadata["step"] ab und ruft prepare_next_tasks auf, um das Aufgaben-dict dieses Schritts zu berechnen.
  • task_states 递归:1199-1227 — baut für jede Aufgabe mit task.name in subgraphs ein checkpoint_ns und ruft rekursiv get_state des Untergraphen auf.
  • apply pending writes:1229-1255 — wendet Schreiboperationen mit NULL_TASK_ID und reguläre pending writes temporär auf die Kanäle an und berechnet dann tasks_with_writes.
  • 组装 StateSnapshot:1257-1266 — setzt Kanalwerte, next-Knoten-tuple, config, metadata, ts, parent_config, tasks, interrupts zum endgültigen StateSnapshot zusammen.
  • get_state:1392-1434 — synchroner Einstieg: löst den checkpointer auf, behandelt das Routing über checkpoint_ns zum Untergraphen und entscheidet über apply_pending_writes.
  • get_state_history:1480 — iteriert über alle historischen Checkpoints derselben thread_id und erzeugt für jeden ein StateSnapshot.
  • _aprepare_state_snapshot:1268 — asynchrone Version; Logik symmetrisch zur synchronen, nur channels_from_checkpoint wird durch achannels_from_checkpoint ersetzt.

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.

python
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:

python
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]),
)

(组装:1257-1266)

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, ob snapshot.next ein leeres Tupel ist, um zwischen «Graph noch nicht gelaufen» und «Graph fertig» zu unterscheiden.
  • get_state ohne checkpointer wirft direkt ValueError(no checkpointer:1401-1402),No checkpointer set — ein rein im Speicher laufender Graph hat keine historischen Snapshots.
  • Wenn checkpoint_ns keinen passenden Untergraphen findet, wird ein Fehler geworfen(subgraph not found:1415-1416),Subgraph {recast} not found; tritt häufig auf, wenn der Namespace fehlerhaft mit NS_SEP (|) oder NS_END (:) zusammengesetzt wurde.
  • apply_pending_writes wird nur aktiviert, wenn checkpoint_id nicht 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 next filtert Aufgaben heraus, die bereits writes haben(next 过滤:1259),Aufgaben mit nicht-leerem t.writes zählen nicht als «next» — da ihre Eingaben in diesem Schritt bereits konsumiert wurden, spiegelt next die «noch nicht gelaufenen Aufgaben» wider.
  • tasks enthält interrupts(interrupts 聚合:1265),die Interrupt-Objekte aller Aufgaben werden zu einem flachen Tupel zusammengefasst; der Benutzer füllt sie in dieser Reihenfolge zurück, wenn er Command(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.