Skip to content

BaseChannel-Abstraktion: Kanal als Zustandsprimitiv

源码版本1.2.9

Verantwortung

In der LangGraph-Laufzeit ist ein Kanal (channel) die kleinste Einheit des Zustands. Eine Pregel-Instanz, die du aus einem StateGraph kompilierst, hält intern nicht direkt das dict des Nutzers, sondern eine Menge von BaseChannel-Instanzen — pro Kanal ein key, und der Kanal ist selbst für Speichern, Zusammenführen von Schreibungen, das Zur-Verfügung-Stellen von Lesezugriffen und für die Serialisierung seines aktuellen Zustands in einen Checkpoint (checkpoint) verantwortlich. BaseChannel ist die abstrakte Basisklasse, die alle Kanal-Klassen teilen (BaseChannel:19).

Sie selbst speichert keine Daten (__slots__ enthält nur key und typ, siehe __slots__:22-26). Sie legt nur vier abstrakte Methoden fest, die Unterklassen implementieren müssen: die beiden Typattribute ValueType / UpdateType sowie die drei Aktionen update / get / from_checkpoint (abstract methods:60-99). Außerdem gibt es fünf Methoden mit Standardimplementierung — checkpoint / is_available / consume / finish / copy —, die entweder auf get() zurückfallen oder False als no-op zurückgeben.

Mit anderen Worten: BaseChannel macht «wie Zustand gespeichert, wie er verändert und wie er persistiert wird» in einem einzigen Interface fest. Die Pregel-Engine muss nur update / get / checkpoint aufrufen, ohne wissen zu müssen, ob dahinter ein LastValue, ein Topic oder ein BinaryOperatorAggregate das Zusammenführen übernimmt. Genau deshalb lassen sich Engine und Zustandsform sauber entkoppeln.

Entwurfsmotivation

Warum Zustand als «Kanal»-Objekt modellieren, statt einfach ein dict zwischen Knoten weiterzureichen?

  • Einheitliche Merge-Semantik: Innerhalb desselben Superstep (superstep) können mehrere Knoten gleichzeitig auf denselben key schreiben. Ein dict kann «Liste anhängen, Skalar überschreiben, Zähler addieren» nicht ausdrücken; ein Kanal bündelt diese Strategie in einem update(values: Sequence[Update]) -> bool, und die Unterklasse entscheidet selbst, wie gefaltet wird.
  • Unveränderliche In-Step-Ansicht: Ein Knoten liest den Snapshot vom Schrittanfang; Werte, die andere Knoten in diesem Schritt schreiben, werden erst im nächsten Schritt sichtbar. get liest nur den aktuellen Wert, update wird nur an der Schrittgrenze von apply_writes aufgerufen (apply_writes:317-323), was Race Conditions konstruktiv ausschließt.
  • Serialisierbar: Jeder Kanal kann per checkpoint() einen serialisierbaren Wert erzeugen (checkpoint:49-58) und sich per from_checkpoint rekonstruieren (from_checkpoint:60-65). Die BaseCheckpointSaver-Familie verlässt sich auf diese beiden Schnittstellen, um den gesamten Thread-Zustand in SQLite/Postgres zu schreiben.
  • Lebenszyklus-Hooks: Die beiden Methoden consume und finish (consume / finish:101-121) geben speziellen Kanälen (etwa Topic(accumulate=False) oder LastValueAfterFinish) die Möglichkeit, am Schrittende oder am Laufende ihren Zustand anzupassen. Standard no-op bedeutet: Für die meisten Kanäle ist diese Semantik nicht erforderlich.
  • Selbstbeschreibende Typen: Die beiden Attribute ValueType / UpdateType (ValueType / UpdateType:28-36) erlauben es Kompilier- und Laufzeit, zurückzufragen, «was dieser Kanal speichert und was er als Schreibung akzeptiert». StateGraph nutzt das in _is_field_channel:1862-1887, um Deklarationen wie Annotated[..., SomeChannel] zu erkennen.

Schlüsseldateien

  • BaseChannel class:19-26 — Abstrakte Basisklasse, generisch über Value / Update / Checkpoint; __slots__ deklariert nur key und typ.
  • ValueType / UpdateType:28-36 — Zwei abstrakte Properties, über die Unterklassen Speicher- und Schreib-Typ deklarieren.
  • checkpoint:49-58 — Standardimplementierung: gibt direkt self.get() zurück, für leere Kanäle MISSING.
  • from_checkpoint:60-65 — Abstrakte Methode: konstruiert aus einem Checkpoint-Wert einen äquivalenten neuen Kanal; muss von Unterklassen implementiert werden.
  • get:69-73 — Abstrakte Lese-Schnittstelle; leere Kanäle werfen EmptyChannelError.
  • is_available:75-85 — Standardimplementierung über get() plus Abfangen von EmptyChannelError; Unterklassen überschreiben das üblicherweise mit einer schnelleren Prüfung.
  • update:89-99 — Abstrakte Schreib-Schnittstelle; Pregel ruft sie am Ende jedes Superstep einmal auf, Reihenfolge beliebig; True bedeutet, der Kanal wurde verändert.
  • consume / finish:101-121 — Lebenszyklus-Hooks, Standard no-op gibt False zurück.
  • apply_writes:315-345 — Die einzige Stelle, an der Pregel die Schreibungen aller Tasks eines Schritts bündelt und auf die Kanäle anwendet; update / consume / finish werden hier zentral gesteuert.
  • _is_field_channel:1862-1887StateGraph verwendet isinstance(item, BaseChannel) zur Kompilierzeit, um Annotated[..., channel] in eine Channel-Instanz zu übersetzen.

Datenfluss

Das Schreiben auf Kanäle geschieht in apply_writes: Pregel bündelt die Schreibungen aller Tasks des Schritts nach Kanal auf und ruft dann update auf (apply_writes:317-323):

python
# Apply writes to channels
updated_channels: set[str] = set()
for chan, vals in pending_writes_by_channel.items():
    if chan in channels:
        if channels[chan].update(vals) and next_version is not None:
            checkpoint["channel_versions"][chan] = next_version
            # unavailable channels can't trigger tasks, so don't add them
            if channels[chan].is_available():
                updated_channels.add(chan)

# Channels that weren't updated in this step are notified of a new step
if bump_step:
    for chan in channels:
        if channels[chan].is_available() and chan not in updated_channels:
            if channels[chan].update(EMPTY_SEQ) and next_version is not None:
                checkpoint["channel_versions"][chan] = next_version
                if channels[chan].is_available():
                    updated_channels.add(chan)

Achte auf den zweiten Block: Kanäle, die in diesem Schritt von keinem Task beschrieben wurden, erhalten ebenfalls einen update(EMPTY_SEQ)-Aufruf. Das ist ein impliziter Teil des BaseChannel-Vertrags — Unterklassen müssen in update leere Sequenzen akzeptieren und in der Regel False («keine Änderung») zurückgeben. Die Semantik «jeden Schritt leeren» von Topic(accumulate=False) baut darauf auf: Ein leeres update veranlasst ihn, den Wert aus dem vorherigen Schritt zu verwerfen, sodass er im nächsten Schritt unsichtbar ist.

Die gesamte Kanal-Lebenszyklus-Position innerhalb eines Pregel-Ticks:

Grenzen und Fehler

  • EmptyChannelError: get() wirft diesen Fehler, wenn der Kanal nie beschrieben wurde (get:70-73). Die Standardimplementierung von is_available leitet die Verfügbarkeit genau über dieses Abfangen ab; Unterklassen mit billigerer Prüfung sollten überschreiben.
  • MISSING-Wächter: checkpoint() gibt langgraph._internal._typing.MISSING statt None zurück, wenn get() ein EmptyChannelError wirft (checkpoint fallback:55-58), damit None als legitimer gespeicherter Wert auftreten kann und nicht mit «leer» verwechselt wird.
  • update gibt False zurück: Das signalisiert, dass der Kanal unverändert blieb. Pregel aktualisiert dann channel_versions nicht, und prepare_next_tasks wird über diesen Kanal keinen Knoten auslösen (update call:319) — Grundlage der Pregel-Konvergenzprüfung.
  • Leere-Sequenz-Aufruf: apply_writes ruft auch für unveränderte Kanäle update(EMPTY_SEQ) auf (EMPTY_SEQ update:329). Der BaseChannel.update-Vertrag ist «leere Eingabe gibt False», aber Topic(accumulate=False) leert sich in diesem Fall und kann True zurückgeben. Unterklassen dürfen also nicht annehmen, «leere Eingabe = nichts tun».
  • consume und finish sind standardmäßig no-op (consume / finish:101-121): Sie geben False zurück, d. h. die meisten Kanäle nehmen am Lebenszyklus nicht teil. apply_writes ruft am Ende von bump_step finish auf (finish call:338), aber nur spezielle Kanäle wie LastValueAfterFinish nutzen das für «erst am Laufende sichtbar».
  • copy läuft standardmäßig über den Checkpoint: BaseChannel.copy ist per Standard self.from_checkpoint(self.checkpoint()) (copy:40-47). Wenn eine Unterklasse teure Checkpoints erzeugt (z. B. tiefe Kopie einer großen list), sollte sie copy überschreiben, um eine billigere Implementierung bereitzustellen — so macht es Topic.copy.

Zusammenfassung

BaseChannel macht «wie Zustand gespeichert, wie zusammengeführt und wie serialisiert wird» in einem Interface fest. Die Pregel-Engine braucht nur die feste Reihenfolge update → is_available → get → checkpoint → finish, um auf jeder beliebigen Kanal-Implementierung zu laufen. Konkrete Implementierungen stehen in /channels/last-value und /channels/topic-binop; wie Pregel Kanäle liest und schreibt, in /pregel/algo; wie die Checkpoint-Saver den Rückgabewert von checkpoint() persistieren, in /checkpoint/base-saver.

Siehe offizielle Dokumentation: LangGraph 文档 · README