Topic und BinaryOperatorAggregate: Akkumulierende Kanäle und Reducer
Verantwortung
LastValue überschreibt nur, aber in echten Agent-Zuständen sollen viele Felder akkumuliert werden: Nachrichtenlisten anhängen, Zähler addieren, Todos zusammenführen. LangGraph bietet dafür zwei Wege: Topic (Topic class:23) ist ein Kanal (channel) im Pub/Sub-Stil für Listen, und BinaryOperatorAggregate (BinaryOperatorAggregate class:65) ist ein generischer Kanal, der einen beliebigen binären Operator (reducer) akzeptiert. Beide erben direkt oder indirekt von BaseChannel (BaseChannel:19).
BinaryOperatorAggregate ist die Implementierung hinter Deklarationen wie Annotated[list, add_messages] — StateGraph ruft zur Kompilierzeit _is_field_binop auf (_is_field_binop:1890-1908), interpretiert das letzte aufrufbare Objekt mit zwei Parametern in den Annotated-Metadaten als reducer und instanziiert BinaryOperatorAggregate(typ, reducer). Bei jedem update wird der neue Wert über operator in den aktuellen Wert eingearbeitet (update:123-144).
Topic wird häufiger für interne Kanäle verwendet: Jede Pregel-Instanz besitzt einen Kanal namens __pregel_tasks vom Typ Topic(Send, accumulate=False) (TASKS channel:809), der die Send-Aufteilung (Fan-out) aufnimmt. Sein Merkmal ist «jeden Schritt leeren» bei accumulate=False — die Send-Liste aus dem vorherigen Schritt wird nach dem Konsumieren durch prepare_next_tasks verworfen und im nächsten Schritt neu gesammelt.
Beide haben gemeinsam, dass update mehrere Werte akzeptiert und nach einer Regel faltet (fold). Darin unterscheiden sie sich grundsätzlich von LastValue, das nur einen Wert akzeptiert.
Entwurfsmotivation
- Reducer ist die generische Schnittstelle für Zustands-Merges: Die Merge-Semantiken verschiedener Felder unterscheiden sich stark —
messageswird nach ID dedupliziert angehängt,counteraddiert,tagswird vereinigt. Statt für jede Variante eine eigene Kanal-Klasse zu bauen, gibt ein generischer Kanal einen binären Operatoroperator: (Value, Value) -> Valuevor und überlässt dem Nutzer die Merge-Regel.add_messages(add_messages:60-66) ist ein solcher reducer: nach Nachrichten-ID zusammenführen, neue ID anhängen, alte ID ersetzen. - Typ-Löschung rückgängig machen:
_strip_extras(_strip_extras:22-28) entfernt die Typing-WrapperAnnotated/Required/NotRequiredund instanziiert danach mittyp()einen Default-Wert. Abstrakte Typen wiecollections.abc.Sequencelassen sich nicht direkt instanziieren, deshalb werden sie intyp fallback:83-88auflist/set/dictabgebildet (typ instantiation:82-92). - Zwei Formen von
Topic.accumulate: Beiaccumulate=Trueist es eine schrittübergreifend akkumulierte Liste —updatehängt neue Werte perextendan (update:77-85); beiaccumulate=Falsewird die Liste zunächst geleert und dann erweitert, sodass nur die Werte des aktuellen Schritts sichtbar bleiben. Letzteres ist genau das, was derTASKS-Kanal braucht: einSendist nur für einen Schritt gültig und wird nach dem Konsumieren ungültig. _flattenunterstützt gemischte Schreibungen:Topic.updateakzeptiert eine gemischte Sequenz ausValue | list[Value](_flatten:15-20), sodass ein Knoten entweder ein einzelnesSendoder[Send, Send]auf einmal schreiben kann — der Kanal flacht das einheitlich ab.Overwrite-Bypass: Ein reducer-Kanal unterstützt auch eine direkte Überschreib-Semantik: DieOverwrite-Datenklasse (Overwrite:980) oder ein dict{"__overwrite__": value}(OVERWRITE key:95) wird von_get_overwriteerkannt (_get_overwrite:31-51) und umgeht inupdateden reducer. Pro Schritt ist nur einOverwriteerlaubt; ein zweites wirftInvalidUpdateError(overwrite dedup:132-138).
Schlüsseldateien
Topic class:23-41— Klassendefinition, generisch überSequence[Value]/Value | list[Value]/list[Value]; der Konstruktor nimmt den Schalteraccumulate._flatten:15-20— Hilfsfunktion: flacht eine gemischteValue | list[Value]-Sequenz zu einemIterator[Value]ab.Topic.copy:56-61— Überschreibtcopy, um direkt eine Kopie der values-Liste wiederzuverwenden, ohne Checkpoint-Serialisierung.Topic.checkpoint / from_checkpoint:63-75—checkpointgibt direkt dieself.values-Liste zurück;from_checkpointist noch mit dem alten tuple-Format kompatibel.Topic.update:77-85— Kernlogik: beiaccumulate=Falseerst leeren, dann extend; beiaccumulate=Trueeinfach extend; leere geflattete Sequenz gibtFalsezurück.Topic.get / is_available:87-94— Leere Liste wirftEmptyChannelError;is_availableist alsbool(self.values)überschrieben.Binop class:65-92— Klassendefinition und__init__: nimmt das binäre Callableoperator, initialisiert den Wert mittyp()._get_overwrite:31-51— Erkennt dreiOverwrite-Darstellungen (Datenklasse, dict mit sentinel-key, nach JSON rekonstruiertes dict)._operators_equal:54-62— Spezielle Lambda-Behandlung: gleichnamige Lambdas gelten als gleich, damit bei jeder Neukompilierung des Graphen nicht jeder Kanal als verändert gilt.Binop.update:123-144— Das Herz des reducer-Kanals: der erste Wert initialisiert, danach wird jeder Wert überoperator(self.value, value)akkumuliert;Overwrite-Bypass ist die Ausnahme._is_field_binop:1890-1908— Erkennt zur Kompilierzeit das binäre Callable in denAnnotated[..., reducer]-Metadaten und instanziiert einBinaryOperatorAggregate.TASKS channel:805-809—Pregel.__init__kodiert hartTopic(Send, accumulate=False)für den__pregel_tasks-Kanal; Nutzer können das nicht überschreiben.add_messages:60-66— Typischer reducer: führt zwei Listen nach Nachrichten-ID zusammen; neue ID anhängen, alte ersetzen.
Datenfluss
BinaryOperatorAggregate.update ist das Herz des reducer-Kanals (update:123-144):
def update(self, values: Sequence[Value]) -> bool:
if not values:
return False
if self.value is MISSING:
self.value = values[0]
values = values[1:]
seen_overwrite: bool = False
for value in values:
is_overwrite, overwrite_value = _get_overwrite(value)
if is_overwrite:
if seen_overwrite:
msg = create_error_message(
message="Can receive only one Overwrite value per super-step.",
error_code=ErrorCode.INVALID_CONCURRENT_GRAPH_UPDATE,
)
raise InvalidUpdateError(msg)
self.value = overwrite_value
seen_overwrite = True
continue
if not seen_overwrite:
self.value = self.operator(self.value, value)
return TrueAn diesem Block sind zwei Zweige wichtig: Wenn der Kanal leer ist, wird der erste Wert direkt als Initialwert gesetzt (ein reducer-Aufruf entfällt, weil der reducer zwei Eingaben braucht); danach wird jeder folgende Wert, sofern er kein Overwrite ist, über self.value = operator(self.value, value) akkumuliert. Sobald ein Overwrite auftritt, werden alle späteren Nicht-Overwrite-Werte verworfen — das stellt sicher, dass die Überschreib-Semantik nicht durch nachfolgende reducer wieder überschrieben wird.
Wie StateGraph aus Annotated[list[AnyMessage], add_messages] ein BinaryOperatorAggregate macht, steht in _is_field_binop:1890-1908:
def _is_field_binop(typ: type[Any]) -> BinaryOperatorAggregate | None:
if hasattr(typ, "__metadata__"):
meta = typ.__metadata__
if len(meta) >= 1 and callable(meta[-1]):
sig = signature(meta[-1])
params = list(sig.parameters.values())
if (
sum(
p.kind in (p.POSITIONAL_ONLY, p.POSITIONAL_OR_KEYWORD)
for p in params
)
== 2
):
return BinaryOperatorAggregate(typ, meta[-1])
else:
raise ValueError(
f"Invalid reducer signature. Expected (a, b) -> c. Got {sig}"
)
return NoneEs validiert mit inspect.signature, dass das letzte Meta-Element ein binäres Callable ist; passt die Signatur nicht, wird direkt ein Fehler geworfen. Deshalb kann eine Funktion wie add_messages(left, right) mit zwei Parametern als reducer dienen, während einparametrige Funktionen abgelehnt werden.
Topic.update nimmt einen anderen, kürzeren Weg (update:77-85):
def update(self, values: Sequence[Value | list[Value]]) -> bool:
updated = False
if not self.accumulate:
updated = bool(self.values)
self.values = list[Value]()
if flat_values := tuple(_flatten(values)):
updated = True
self.values.extend(flat_values)
return updatedBeachte den Rückgabewert bei accumulate=False: Wenn der alte Wert nicht leer war, wird True zurückgegeben, selbst wenn die neue Sequenz leer ist (weil das Leeren selbst eine Änderung ist), und Pregel aktualisiert channel_versions. Das ist das Gegenteil von LastValue, das bei leerer Sequenz False zurückgibt — unterschiedliche Kanäle haben unterschiedliche Semantik für «leeres Schreiben».
Position eines reducer-Kanals innerhalb eines Pregel-Superstep:
Grenzen und Fehler
- Drei
Overwrite-Erkennungsformen parallel:overwrite forms:44-50akzeptiert gleichzeitig dieOverwrite-Datenklasse, das dict{"__overwrite__": value}und das nach JSON rekonstruierte dict{"type": "__overwrite__", "value": ...}. Der dritte Zweig stellt sicher, dass die Semantik nach einer orjson-Roundtrip über den LangGraph API server erhalten bleibt. - Pro Schritt nur ein
Overwrite:overwrite dedup:132-138wirft beim zweiten Auftreten direktInvalidUpdateError, weil «Überschreiben» mehrfach im selben Schritt widersprüchlich ist. - Nicht-Overwrite-Werte nach
Overwritewerden still verworfen:discard after overwrite:142-143—if not seen_overwriteentscheidet, ob der reducer aufgerufen wird. Sobaldseen_overwrite=True, werden spätere Werte weder in den reducer eingespeist noch als Fehler geworfen. Das ist beabsichtigt: «Überschreiben» heißt, dass das bisherige Merge-Ergebnis nicht erhalten bleiben soll und auch keine nachfolgenden Werte nötig sind. - Vergleich von Lambda-Redcern auf Gleichheit:
_operators_equal:54-62behandelt alle Lambdas als gleich, weil ein Python-Lambda bei jedem<lambda>ein neues Objekt ist. Ein identitätsbasierter Vergleich würde__eq__nie erfüllen und den Graph-Cache stören. Topic(accumulate=False)gibt bei leerer SchreibungTruezurück:clear returns True:79-80— solange der alte Wert nicht leer war, gilt das Leeren als Änderung. Das veranlasstapply_writes,channel_versionszu aktualisieren und Knoten auszulösen, die von diesem Kanal abhängen. WerTopicselbst nutzt, sollte das wissen: Es behandelt leere Schreibungen nicht so stillschweigend wieLastValue.Topicist mit alten Checkpoints kompatibel:tuple compatibility:70-74prüftisinstance(checkpoint, tuple); alte Checkpoints speicherten(typ, values)als tuple, neuere direkt alslist. Dieser Code bricht alte Archive nicht.- Reducer-Signatur strikt binär:
reducer signature check:1894-1907zählt mitinspect.signaturedie Positions-Parameter; bei nicht genau zwei wird ein Fehler geworfen. Funktionen mit*argswerden nicht akzeptiert — die Regel «reducer muss strikt(a, b) -> csein» ist verbindlich. add_messagesist eine Factory:_add_messages_wrapper:41-57dekoriert mit@_add_messages_wrapper. Die eigentliche Funktion lässt sich alsadd_messages(left, right)direkt aufrufen oder alsadd_messages(format="langchain-openai")eine partial zurückgeben. Dieses Muster ist eine häufige Erweiterungsstelle des reducer-Musters — dem reducer Laufzeit-Argumente mitgeben.
Zusammenfassung
BinaryOperatorAggregate ist die einheitliche Implementierung für «Felder mit reducer»; Topic ist die zwei-Formen-Liste: schrittübergreifend akkumulieren oder jeden Schritt leeren. Zusammen mit /channels/last-value bilden sie das Trio des LangGraph-Kanalsystems; die Schnittstelle kommt aus /channels/base-channel. Ein reducer wie add_messages wird in /graph/state-graph per Annotated[..., reducer] deklariert und zur Laufzeit über apply_writes in /pregel/algo in update gerufen. Topic(accumulate=False) ist als __pregel_tasks-Kanal der Träger der Send-Aufteilung, siehe /subgraph/send-command.
Siehe offizielle Dokumentation: LangGraph 文档 · README。