Skip to content

Topic und BinaryOperatorAggregate: Akkumulierende Kanäle und Reducer

源码版本1.2.9

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 — messages wird nach ID dedupliziert angehängt, counter addiert, tags wird vereinigt. Statt für jede Variante eine eigene Kanal-Klasse zu bauen, gibt ein generischer Kanal einen binären Operator operator: (Value, Value) -> Value vor 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-Wrapper Annotated / Required / NotRequired und instanziiert danach mit typ() einen Default-Wert. Abstrakte Typen wie collections.abc.Sequence lassen sich nicht direkt instanziieren, deshalb werden sie in typ fallback:83-88 auf list / set / dict abgebildet (typ instantiation:82-92).
  • Zwei Formen von Topic.accumulate: Bei accumulate=True ist es eine schrittübergreifend akkumulierte Liste — update hängt neue Werte per extend an (update:77-85); bei accumulate=False wird die Liste zunächst geleert und dann erweitert, sodass nur die Werte des aktuellen Schritts sichtbar bleiben. Letzteres ist genau das, was der TASKS-Kanal braucht: ein Send ist nur für einen Schritt gültig und wird nach dem Konsumieren ungültig.
  • _flatten unterstützt gemischte Schreibungen: Topic.update akzeptiert eine gemischte Sequenz aus Value | list[Value] (_flatten:15-20), sodass ein Knoten entweder ein einzelnes Send oder [Send, Send] auf einmal schreiben kann — der Kanal flacht das einheitlich ab.
  • Overwrite-Bypass: Ein reducer-Kanal unterstützt auch eine direkte Überschreib-Semantik: Die Overwrite-Datenklasse (Overwrite:980) oder ein dict {"__overwrite__": value} (OVERWRITE key:95) wird von _get_overwrite erkannt (_get_overwrite:31-51) und umgeht in update den reducer. Pro Schritt ist nur ein Overwrite erlaubt; ein zweites wirft InvalidUpdateError (overwrite dedup:132-138).

Schlüsseldateien

  • Topic class:23-41 — Klassendefinition, generisch über Sequence[Value] / Value | list[Value] / list[Value]; der Konstruktor nimmt den Schalter accumulate.
  • _flatten:15-20 — Hilfsfunktion: flacht eine gemischte Value | list[Value]-Sequenz zu einem Iterator[Value] ab.
  • Topic.copy:56-61 — Überschreibt copy, um direkt eine Kopie der values-Liste wiederzuverwenden, ohne Checkpoint-Serialisierung.
  • Topic.checkpoint / from_checkpoint:63-75checkpoint gibt direkt die self.values-Liste zurück; from_checkpoint ist noch mit dem alten tuple-Format kompatibel.
  • Topic.update:77-85 — Kernlogik: bei accumulate=False erst leeren, dann extend; bei accumulate=True einfach extend; leere geflattete Sequenz gibt False zurück.
  • Topic.get / is_available:87-94 — Leere Liste wirft EmptyChannelError; is_available ist als bool(self.values) überschrieben.
  • Binop class:65-92 — Klassendefinition und __init__: nimmt das binäre Callable operator, initialisiert den Wert mit typ().
  • _get_overwrite:31-51 — Erkennt drei Overwrite-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 über operator(self.value, value) akkumuliert; Overwrite-Bypass ist die Ausnahme.
  • _is_field_binop:1890-1908 — Erkennt zur Kompilierzeit das binäre Callable in den Annotated[..., reducer]-Metadaten und instanziiert ein BinaryOperatorAggregate.
  • TASKS channel:805-809Pregel.__init__ kodiert hart Topic(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):

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

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

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

Es 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):

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

Beachte 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-50 akzeptiert gleichzeitig die Overwrite-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-138 wirft beim zweiten Auftreten direkt InvalidUpdateError, weil «Überschreiben» mehrfach im selben Schritt widersprüchlich ist.
  • Nicht-Overwrite-Werte nach Overwrite werden still verworfen: discard after overwrite:142-143if not seen_overwrite entscheidet, ob der reducer aufgerufen wird. Sobald seen_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-62 behandelt 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 Schreibung True zurück: clear returns True:79-80 — solange der alte Wert nicht leer war, gilt das Leeren als Änderung. Das veranlasst apply_writes, channel_versions zu aktualisieren und Knoten auszulösen, die von diesem Kanal abhängen. Wer Topic selbst nutzt, sollte das wissen: Es behandelt leere Schreibungen nicht so stillschweigend wie LastValue.
  • Topic ist mit alten Checkpoints kompatibel: tuple compatibility:70-74 prüft isinstance(checkpoint, tuple); alte Checkpoints speicherten (typ, values) als tuple, neuere direkt als list. Dieser Code bricht alte Archive nicht.
  • Reducer-Signatur strikt binär: reducer signature check:1894-1907 zählt mit inspect.signature die Positions-Parameter; bei nicht genau zwei wird ein Fehler geworfen. Funktionen mit *args werden nicht akzeptiert — die Regel «reducer muss strikt (a, b) -> c sein» ist verbindlich.
  • add_messages ist eine Factory: _add_messages_wrapper:41-57 dekoriert mit @_add_messages_wrapper. Die eigentliche Funktion lässt sich als add_messages(left, right) direkt aufrufen oder als add_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