BaseChannel-Abstraktion: Kanal als Zustandsprimitiv
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.
getliest nur den aktuellen Wert,updatewird nur an der Schrittgrenze vonapply_writesaufgerufen (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 perfrom_checkpointrekonstruieren (from_checkpoint:60-65). DieBaseCheckpointSaver-Familie verlässt sich auf diese beiden Schnittstellen, um den gesamten Thread-Zustand in SQLite/Postgres zu schreiben. - Lebenszyklus-Hooks: Die beiden Methoden
consumeundfinish(consume / finish:101-121) geben speziellen Kanälen (etwaTopic(accumulate=False)oderLastValueAfterFinish) 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».StateGraphnutzt das in_is_field_channel:1862-1887, um Deklarationen wieAnnotated[..., SomeChannel]zu erkennen.
Schlüsseldateien
BaseChannel class:19-26— Abstrakte Basisklasse, generisch überValue / Update / Checkpoint;__slots__deklariert nurkeyundtyp.ValueType / UpdateType:28-36— Zwei abstrakte Properties, über die Unterklassen Speicher- und Schreib-Typ deklarieren.checkpoint:49-58— Standardimplementierung: gibt direktself.get()zurück, für leere KanäleMISSING.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 werfenEmptyChannelError.is_available:75-85— Standardimplementierung überget()plus Abfangen vonEmptyChannelError; 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;Truebedeutet, der Kanal wurde verändert.consume / finish:101-121— Lebenszyklus-Hooks, Standard no-op gibtFalsezurü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/finishwerden hier zentral gesteuert._is_field_channel:1862-1887—StateGraphverwendetisinstance(item, BaseChannel)zur Kompilierzeit, umAnnotated[..., 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):
# 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 vonis_availableleitet die Verfügbarkeit genau über dieses Abfangen ab; Unterklassen mit billigerer Prüfung sollten überschreiben.MISSING-Wächter:checkpoint()gibtlanggraph._internal._typing.MISSINGstattNonezurück, wennget()einEmptyChannelErrorwirft (checkpoint fallback:55-58), damitNoneals legitimer gespeicherter Wert auftreten kann und nicht mit «leer» verwechselt wird.updategibtFalsezurück: Das signalisiert, dass der Kanal unverändert blieb. Pregel aktualisiert dannchannel_versionsnicht, undprepare_next_taskswird über diesen Kanal keinen Knoten auslösen (update call:319) — Grundlage der Pregel-Konvergenzprüfung.- Leere-Sequenz-Aufruf:
apply_writesruft auch für unveränderte Kanäleupdate(EMPTY_SEQ)auf (EMPTY_SEQ update:329). DerBaseChannel.update-Vertrag ist «leere Eingabe gibtFalse», aberTopic(accumulate=False)leert sich in diesem Fall und kannTruezurückgeben. Unterklassen dürfen also nicht annehmen, «leere Eingabe = nichts tun». consumeundfinishsind standardmäßig no-op (consume / finish:101-121): Sie gebenFalsezurück, d. h. die meisten Kanäle nehmen am Lebenszyklus nicht teil.apply_writesruft am Ende vonbump_stepfinishauf (finish call:338), aber nur spezielle Kanäle wieLastValueAfterFinishnutzen das für «erst am Laufende sichtbar».copyläuft standardmäßig über den Checkpoint:BaseChannel.copyist per Standardself.from_checkpoint(self.checkpoint())(copy:40-47). Wenn eine Unterklasse teure Checkpoints erzeugt (z. B. tiefe Kopie einer großen list), sollte siecopyüberschreiben, um eine billigere Implementierung bereitzustellen — so macht esTopic.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。