Topic et BinaryOperatorAggregate : canaux accumulants et réducteurs
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 —
messagesdoit dédupliquer par ID,counteradditionne,tagsfait 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 binaireoperator: (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 wrappersAnnotated/Required/NotRequired, puis utilisetyp()pour instancier une valeur par défaut. Les types abstraits commecollections.abc.Sequencene sont pas directement instanciables, donc on les remonte verslist/set/dictdanstyp fallback:83-88(typ instantiation:82-92). - Double forme via
Topic.accumulate: avecaccumulate=True, c'est une liste qui s'accumule entre superpas —updateyextendles nouvelles valeurs (update:77-85) ; avecaccumulate=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 canalTASKSa besoin : unSendn'est valable qu'un superpas, il devient obsolète une fois consommé. _flattengère les écritures mixtes :Topic.updateaccepte une séquence mixteValue | list[Value](_flatten:15-20), donc un nœud peut écrire un seulSendou[Send, Send], le canal aplatit tout.- Bypass
Overwrite: les canaux à reducer supportent aussi une sémantique « écrasement direct » : le dataclassOverwrite(Overwrite:980) ou le dictionnaire{"__overwrite__": value}(OVERWRITE key:95) est reconnu par_get_overwrite(_get_overwrite:31-51) et bypass le reducer dansupdatepour une affectation directe. Un seulOverwriteest autorisé par superpas ; un second lèveInvalidUpdateError(overwrite dedup:132-138).
Fichiers clés
Topic class:23-41— définition, paramètres génériquesSequence[Value]/Value | list[Value]/list[Value], constructeur acceptant un switchaccumulate._flatten:15-20— utilitaire qui aplatit une séquence mixteValue | list[Value]enIterator[Value].Topic.copy:56-61— surchargecopypour réutiliser une copie de la listevalues, sans passer par la sérialisation checkpoint.Topic.checkpoint / from_checkpoint:63-75— checkpoint renvoie directementself.values;from_checkpointreste compatible avec l'ancien format tuple.Topic.update:77-85— logique centrale :accumulate=Falseefface puis extend,accumulate=Trueextend directement ; séquence aplatie vide renvoieFalse.Topic.get / is_available:87-94— liste vide lèveEmptyChannelError;is_availablesurchargée enbool(self.values).Binop class:65-92— définition +__init__: accepte un binaireoperatorcallable, initialise la valeur viatyp()._get_overwrite:31-51— reconnaît trois formes deOverwrite(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 paroperator(self.value, value), sauf bypassOverwrite._is_field_binop:1890-1908— à la compilation, reconnaît le callable à deux paramètres dans les métadonnéesAnnotated[..., reducer]et instancieBinaryOperatorAggregate.TASKS channel:805-809—Pregel.__init__code en dur unTopic(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) :
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 :
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 Noneinspect.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) :
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 updatedNotez 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
Overwritecoexistent :overwrite forms:44-50accepte simultanément le dataclassOverwrite, 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
Overwritepar superpas :overwrite dedup:132-138lèveInvalidUpdateErrordès la seconde occurrence, car deux écrasements dans le même superpas seraient contradictoires. - Les valeurs non-overwrite suivant un
Overwritesont silencieusement jetées :discard after overwrite:142-143if not seen_overwritedécide d'appeler le reducer ; dès queseen_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-62traite 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)renvoieTruesur écriture vide :clear returns True:79-80— tant que l'ancienne valeur est non vide, l'effacement compte comme modification. Cela pousseapply_writesà mettre à jourchannel_versionset donc à déclencher les nœuds dépendant de ce canal. Si vous utilisezTopicvous-même, soyez-en conscient : contrairement àLastValue, il ne traite pas silencieusement l'écriture vide.- Compatibilité des anciens checkpoints
Topic:tuple compatibility:70-74vérifieisinstance(checkpoint, tuple); l'ancien format stockait un tuple(typ, values), le nouveau stocke directement unelist. Ce bloc évite de casser les anciens checkpoints. - Signature reducer strictement à deux paramètres :
reducer signature check:1894-1907utiliseinspect.signaturepour compter les paramètres positionnels ; tout écart lève une erreur. Une fonction à*argsn'est pas acceptée — « reducer doit être strictement(a, b) -> c». add_messagesest une factory :_add_messages_wrapper:41-57utilise@_add_messages_wrapper, la fonction peut être appelée directement commeadd_messages(left, right)ou commeadd_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