Skip to content

Topic et BinaryOperatorAggregate : canaux accumulants et réducteurs

源码版本1.2.9

Responsabilités

LastValue ne fait qu'écraser, mais dans un état d'agent réel beaucoup de champs doivent s'accumuler : liste de messages à étendre, compteur à additionner, tâches à fusionner. LangGraph offre deux voies pour ce besoin : Topic (Topic class:23), un canal de liste de style pub/sub, et BinaryOperatorAggregate (BinaryOperatorAggregate class:65), un canal générique acceptant n'importe quel opérateur binaire (reducer). Tous deux héritent, directement ou indirectement, de BaseChannel (BaseChannel:19).

BinaryOperatorAggregate est l'implémentation derrière les déclarations Annotated[list, add_messages] écrites par l'utilisateur — à la compilation, StateGraph appelle _is_field_binop (_is_field_binop:1890-1908), prend le dernier callable à deux paramètres dans les métadonnées Annotated comme reducer, et instancie BinaryOperatorAggregate(typ, reducer). Chaque update fusionne la nouvelle valeur dans la valeur courante via operator (update:123-144).

Topic est davantage utilisé pour les canaux internes : chaque instance Pregel embarque un canal Topic(Send, accumulate=False) nommé __pregel_tasks (TASKS channel:809) destiné à recevoir le fan-out de Send. Sa particularité est qu'en mode accumulate=False il s'efface à chaque superpas : la liste de Send du superpas précédent est consommée par prepare_next_tasks puis jetée, et le superpas suivant recommence à zéro.

Le dénominateur commun est « update accepte plusieurs valeurs et les fold selon une règle », ce qui les distingue fondamentalement de LastValue — ce dernier n'accepte qu'une seule valeur.

Motivation de conception

  • Le reducer est l'interface générique de fusion d'état : les sémantiques de fusion varient selon le champ — messages doit dédupliquer par ID, counter additionne, tags fait l'union. Plutôt que d'écrire une classe de canal par cas, on fournit un canal générique qui accepte un opérateur binaire operator: (Value, Value) -> Value, laissant l'utilisateur définir sa règle. add_messages (add_messages:60-66) est un tel reducer : il fusionne par ID de message, avec append des nouveaux ID et remplacement des ID existants.
  • Restauration après effacement de type : _strip_extras (_strip_extras:22-28) retire les wrappers Annotated / Required / NotRequired, puis utilise typ() pour instancier une valeur par défaut. Les types abstraits comme collections.abc.Sequence ne sont pas directement instanciables, donc on les remonte vers list / set / dict dans typ fallback:83-88 (typ instantiation:82-92).
  • Double forme via Topic.accumulate : avec accumulate=True, c'est une liste qui s'accumule entre superpas — update y extend les nouvelles valeurs (update:77-85) ; avec accumulate=False, les anciennes valeurs sont d'abord effacées avant l'extend, ne laissant que « ce qui a été produit ce superpas ». Cette dernière forme est exactement ce dont le canal TASKS a besoin : un Send n'est valable qu'un superpas, il devient obsolète une fois consommé.
  • _flatten gère les écritures mixtes : Topic.update accepte une séquence mixte Value | list[Value] (_flatten:15-20), donc un nœud peut écrire un seul Send ou [Send, Send], le canal aplatit tout.
  • Bypass Overwrite : les canaux à reducer supportent aussi une sémantique « écrasement direct » : le dataclass Overwrite (Overwrite:980) ou le dictionnaire {"__overwrite__": value} (OVERWRITE key:95) est reconnu par _get_overwrite (_get_overwrite:31-51) et bypass le reducer dans update pour une affectation directe. Un seul Overwrite est autorisé par superpas ; un second lève InvalidUpdateError (overwrite dedup:132-138).

Fichiers clés

  • Topic class:23-41 — définition, paramètres génériques Sequence[Value] / Value | list[Value] / list[Value], constructeur acceptant un switch accumulate.
  • _flatten:15-20 — utilitaire qui aplatit une séquence mixte Value | list[Value] en Iterator[Value].
  • Topic.copy:56-61 — surcharge copy pour réutiliser une copie de la liste values, sans passer par la sérialisation checkpoint.
  • Topic.checkpoint / from_checkpoint:63-75 — checkpoint renvoie directement self.values ; from_checkpoint reste compatible avec l'ancien format tuple.
  • Topic.update:77-85 — logique centrale : accumulate=False efface puis extend, accumulate=True extend directement ; séquence aplatie vide renvoie False.
  • Topic.get / is_available:87-94 — liste vide lève EmptyChannelError ; is_available surchargée en bool(self.values).
  • Binop class:65-92 — définition + __init__ : accepte un binaire operator callable, initialise la valeur via typ().
  • _get_overwrite:31-51 — reconnaît trois formes de Overwrite (dataclass, dict à clé sentinel, dict restauré depuis JSON).
  • _operators_equal:54-62 — traitement spécial des lambdas : deux lambdas de même nom sont considérées égales, pour éviter qu'une recompilation du graphe ne fasse basculer le canal en « modifié ».
  • Binop.update:123-144 — cœur du reducer : la première valeur sert à initialiser, les suivantes passent par operator(self.value, value), sauf bypass Overwrite.
  • _is_field_binop:1890-1908 — à la compilation, reconnaît le callable à deux paramètres dans les métadonnées Annotated[..., reducer] et instancie BinaryOperatorAggregate.
  • TASKS channel:805-809Pregel.__init__ code en dur un Topic(Send, accumulate=False) pour __pregel_tasks, sans permettre à l'utilisateur de le surcharger.
  • add_messages:60-66 — reducer typique : fusionne deux listes par ID de message, avec append des nouveaux ID et remplacement des anciens.

Flux de données

BinaryOperatorAggregate.update est le cœur du canal à reducer (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

À noter dans cette lecture : la première valeur, quand le canal est vide, devient la valeur initiale (on économise un appel au reducer, qui nécessite deux entrées). Ensuite, chaque valeur qui n'est pas un Overwrite passe par self.value = operator(self.value, value). Si un Overwrite apparaît, toutes les valeurs suivantes non-overwrite sont silencieusement jetées — cela garantit que la sémantique d'écrasement ne peut pas être inversée par un reducer ultérieur.

Comment StateGraph compile-t-il Annotated[list[AnyMessage], add_messages] en BinaryOperatorAggregate ? Voir _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

inspect.signature vérifie que le dernier élément de meta est bien un callable à deux paramètres ; toute signature non conforme lève une erreur. C'est pourquoi add_messages(left, right) à deux paramètres peut servir de reducer, mais qu'une fonction à un seul paramètre serait refusée.

Topic.update suit un autre chemin, plus court (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

Notez le retour en mode accumulate=False : si l'ancienne valeur n'était pas vide, on renvoie True même avec une séquence vide en entrée (parce que l'effacement est lui-même une modification), et Pregel met à jour channel_versions. À l'inverse de LastValue qui renvoie False sur séquence vide, cela reflète des sémantiques distinctes de « l'écriture vide » selon le canal.

Position du canal à reducer dans un superpas Pregel :

Limites et échecs

  • Trois formes de reconnaissance Overwrite coexistent : overwrite forms:44-50 accepte simultanément le dataclass Overwrite, le dictionnaire {"__overwrite__": value}, et le dictionnaire {"type": "__overwrite__", "value": ...} restauré depuis JSON. Cette troisième branche sert à préserver la sémantique après un aller-retour orjson via le LangGraph API server.
  • Un seul Overwrite par superpas : overwrite dedup:132-138 lève InvalidUpdateError dès la seconde occurrence, car deux écrasements dans le même superpas seraient contradictoires.
  • Les valeurs non-overwrite suivant un Overwrite sont silencieusement jetées : discard after overwrite:142-143 if not seen_overwrite décide d'appeler le reducer ; dès que seen_overwrite=True, les valeurs suivantes ne vont ni dans le reducer ni ne lèvent d'erreur. C'est intentionnel : l'« écrasement » signifie qu'on ne veut ni conserver le résultat précédent ni traiter les valeurs suivantes.
  • Égalité des lambdas reducer : _operators_equal:54-62 traite toutes les lambdas comme égales, car en Python chaque <lambda> est un nouvel objet ; une comparaison par identité rendrait __eq__ toujours faux et casserait le cache du graphe.
  • Topic(accumulate=False) renvoie True sur écriture vide : clear returns True:79-80 — tant que l'ancienne valeur est non vide, l'effacement compte comme modification. Cela pousse apply_writes à mettre à jour channel_versions et donc à déclencher les nœuds dépendant de ce canal. Si vous utilisez Topic vous-même, soyez-en conscient : contrairement à LastValue, il ne traite pas silencieusement l'écriture vide.
  • Compatibilité des anciens checkpoints Topic : tuple compatibility:70-74 vérifie isinstance(checkpoint, tuple) ; l'ancien format stockait un tuple (typ, values), le nouveau stocke directement une list. Ce bloc évite de casser les anciens checkpoints.
  • Signature reducer strictement à deux paramètres : reducer signature check:1894-1907 utilise inspect.signature pour compter les paramètres positionnels ; tout écart lève une erreur. Une fonction à *args n'est pas acceptée — « reducer doit être strictement (a, b) -> c ».
  • add_messages est une factory : _add_messages_wrapper:41-57 utilise @_add_messages_wrapper, la fonction peut être appelée directement comme add_messages(left, right) ou comme add_messages(format="langchain-openai") pour renvoyer un partial. C'est un point d'extension courant du motif reducer — passer des paramètres runtime au reducer.

Résumé

BinaryOperatorAggregate est l'implémentation unifiée des « champs avec reducer », Topic est la forme « liste accumulante » en deux variantes : accumulation inter-superpas ou effacement à chaque superpas. Avec /channels/last-value, ils forment le trio du système de canaux LangGraph, dont l'interface vient de /channels/base-channel. Un reducer comme add_messages est déclaré via Annotated[..., reducer] dans /graph/state-graph et appelé au runtime par le update de apply_writes dans /pregel/algo. Topic(accumulate=False) est le canal __pregel_tasks qui porte le fan-out Send, voir /subgraph/send-command.

Voir la documentation officielle : documentation LangGraph · README