Topic y BinaryOperatorAggregate: canal acumulativo y reductor
Responsabilidades
LastValue solo sobrescribe, pero en el estado real de un agente muchos campos se acumulan: lista de mensajes se appenda, contadores se suman, listas de pendientes se fusionan. LangGraph ofrece dos vías para este tipo de necesidad: Topic (Topic class:23) es un canal de lista estilo pub/sub; BinaryOperatorAggregate (BinaryOperatorAggregate class:65) es un canal genérico que acepta cualquier operador binario (reducer). Ambos heredan, directa o indirectamente, de BaseChannel (BaseChannel:19).
BinaryOperatorAggregate es la implementación detrás de declaraciones tipo Annotated[list, add_messages]: al compilar, StateGraph invoca _is_field_binop (_is_field_binop:1890-1908), toma el último callable binario de los metadatos de Annotated como reducer y lo instancia como BinaryOperatorAggregate(typ, reducer). Cada update fusiona los nuevos valores al valor actual con operator (update:123-144).
Topic se usa más en canales internos: cada instancia Pregel lleva un canal Topic(Send, accumulate=False) de nombre __pregel_tasks (TASKS channel:809) para albergar los fan-out Send. Su característica es que con accumulate=False se vacía en cada paso — la lista de Send del paso anterior se descarta tras ser consumida por prepare_next_tasks, y el siguiente paso la recolecta de nuevo.
Ambos comparten «update acepta múltiples valores y los fold según alguna regla», su diferencia radical con LastValue — este último solo acepta un valor.
Motivación de diseño
- reducer es la interfaz universal de fusión de estado: las semánticas de fusión de distintos campos varían enormen —
messagesse appenda deduplicando por ID,counterse suma,tagsse une por unión. En vez de fabricar una clase de canal para cada caso, mejor un canal genérico que acepte un operador binariooperator: (Value, Value) -> Valuey deje al usuario definir la regla.add_messages(add_messages:60-66) es un reducer así: fusiona por ID de mensaje, los IDs nuevos se appendan, los IDs antiguos se reemplazan. - Reconstrucción tras el borrado de tipos:
_strip_extras(_strip_extras:22-28) despoja los envoltorios typing deAnnotated/Required/NotRequiredy luego instanciatyp()para obtener el valor por defecto. Tipos abstractos comocollections.abc.Sequenceno se pueden instanciar directamente, así que entyp fallback:83-88se mapean alist/set/dict(typ instantiation:82-92). Topic.accumulateen doble forma: conaccumulate=Truees una lista acumulativa entre pasos —updatehaceextendde los nuevos valores (update:77-85); conaccumulate=Falseprimero vacía los valores viejos y luego hace extend, convirtiéndose en «ver solo los valores producidos en este paso». Esto último es exactamente lo que necesita el canalTASKS:Sendes válido dentro de un paso y caduca tras ser consumido._flattenadmite escrituras mixtas:Topic.updateacepta una secuencia mixtaValue | list[Value](_flatten:15-20), de modo que un nodo puede escribir unSendo[Send, Send]en una llamada, y el canal lo aplana uniformemente.- Bypass
Overwrite: el canal reducer admite además una semántica de «sobrescritura directa»: el dataclassOverwrite(Overwrite:980) o el dict{"__overwrite__": value}(OVERWRITE key:95) son reconocidos por_get_overwrite(_get_overwrite:31-51) y, dentro deupdate, saltan el reducer y asignan directamente. Solo se permite unOverwritepor paso; un segundo lanzaInvalidUpdateError(overwrite dedup:132-138).
Archivos clave
Topic class:23-41— definición de la clase; parámetros genéricosSequence[Value]/Value | list[Value]/list[Value]; el constructor acepta el flagaccumulate._flatten:15-20— función utilitaria: aplana una secuencia mixtaValue | list[Value]aIterator[Value].Topic.copy:56-61— sobreescribecopypara reutilizar directamente una copia de la lista de valores, sin pasar por la serialización del checkpoint.Topic.checkpoint / from_checkpoint:63-75— checkpoint devuelve directamenteself.values;from_checkpointsigue siendo compatible con el formato antiguo en forma de tuple.Topic.update:77-85— lógica central: conaccumulate=Falseprimero vacía y luego extiende; conaccumulate=Trueextiende directamente; una secuencia flat vacía devuelveFalse.Topic.get / is_available:87-94— lista vacía lanzaEmptyChannelError;is_availablese sobreescribe comobool(self.values).Binop class:65-92— definición de la clase +__init__: acepta el callable binariooperatore inicializa el valor contyp()._get_overwrite:31-51— reconoce tres representaciones deOverwrite(dataclass, dict con clave centinela, dict restaurado desde JSON)._operators_equal:54-62— tratamiento especial para lambda: lambdas del mismo nombre se consideran iguales, para evitar que cada recompilación del grafo marque el canal como cambiado.Binop.update:123-144— la verdadera entrada del reducer: el primer valor inicializa, los siguientes invocanoperator(self.value, value)uno a uno; el bypassOverwritees la excepción._is_field_binop:1890-1908— en compilación reconoce el callable binario en los metadatos deAnnotated[..., reducer]y lo instancia comoBinaryOperatorAggregate.TASKS channel:805-809—Pregel.__init__fija el canal__pregel_taskscomoTopic(Send, accumulate=False), sin permitir que el usuario lo sobrescriba.add_messages:60-66— reducer típico: fusiona dos listas por ID de mensaje, los IDs nuevos se appendan, los antiguos se reemplazan.
Flujo de datos
BinaryOperatorAggregate.update es el corazón del 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 TrueAl leer este fragmento hay que notar dos ramas: si el canal está vacío, el primer valor se usa como inicial (saltando una llamada al reducer, que necesita dos entradas); después, cualquier valor que no sea Overwrite ejecuta self.value = operator(self.value, value). Si aparece un Overwrite, todos los valores posteriores no-overwrite se descartan — garantiza que la semántica de «sobrescritura» no sea revertida por el reducer.
Cómo StateGraph compila Annotated[list[AnyMessage], add_messages] a BinaryOperatorAggregate, ver _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 NoneUsa inspect.signature para validar que el último meta sea un callable binario; si la firma no coincide, lanza error. Por eso add_messages(left, right) (función de dos parámetros) puede usarse como reducer, mientras que una función de un parámetro se rechaza.
Topic.update va por otra ruta, con una lógica más corta (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 updatedNótese el valor de retorno con accumulate=False: si el valor viejo no estaba vacío, devuelve True incluso cuando la secuencia nueva es vacía (porque vaciar es en sí un cambio), y Pregel actualizará channel_versions. Es lo opuesto a LastValue, que devuelve False ante secuencia vacía; refleja la semántica distinta de cada canal ante la «escritura vacía».
La posición del canal reducer dentro de un superpaso de Pregel:
Límites y fallos
- Tres formas reconocibles de
Overwritecoexisten:overwrite forms:44-50acepta simultáneamente el dataclassOverwrite, el dict{"__overwrite__": value}y el dict restaurado desde JSON{"type": "__overwrite__", "value": ...}. La tercera rama existe para que la semántica sobreviva a una ronda de serialización orjson a través del LangGraph API server sin pérdidas. - Solo un
Overwritepor paso:overwrite dedup:132-138lanzaInvalidUpdateErrorante una segunda ocurrencia, porque la semántica de «sobrescritura» es contradictoria si aparece varias veces en un mismo paso. - Los valores no-overwrite tras
Overwritese descartan en silencio:discard after overwrite:142-143;if not seen_overwritedecide si invocar el reducer; una vezseen_overwrite=True, los valores posteriores ni entran al reducer ni generan error. Es intencional — «sobrescritura» significa que no se quiere conservar el resultado previo fusionado ni los valores siguientes. - Comparación de equivalencia de lambdas reducer:
_operators_equal:54-62trata todas las lambdas como iguales, porque en Python cada<lambda>es un objeto nuevo; comparar por identidad haría que__eq__fuera siempre falso y perjudicaría a la caché del grafo. Topic(accumulate=False)devuelveTrueante escritura vacía:clear returns True:79-80— siempre que el valor viejo no estuviera vacío, vaciar cuenta como cambio. Esto hará queapply_writesactualicechannel_versionsy dispare los nodos que dependen de ese canal. Si usasTopicpor tu cuenta, sé consciente: no maneja la escritura vacía en silencio comoLastValue.Topiccompatible con checkpoints antiguos:tuple compatibility:70-74compruebaisinstance(checkpoint, tuple); en versiones antiguas el checkpoint almacenaba un tuple(typ, values)y en la nueva versión se almacena directamente unalist. Este bloque existe para no romper los archivos antiguos.- La firma del reducer se fuerza a dos parámetros:
reducer signature check:1894-1907usainspect.signaturepara contar los parámetros posicionales; si no son exactamente dos, lanza error. Funciones con*argsu otros variadic no se aceptan — fija «el reducer debe ser estrictamente(a, b) -> c». add_messageses una factoría:_add_messages_wrapper:41-57la decora con@_add_messages_wrapper; la función real puede invocarse directamente comoadd_messages(left, right)o comoadd_messages(format="langchain-openai")para devolver un partial. Este patrón es un punto de extensión habitual del modelo reducer — añadir argumentos en tiempo de ejecución al reducer.
Resumen
BinaryOperatorAggregate es la implementación unificada de «campos con reducer»; Topic es la doble forma del «acumulador de lista»: acumulativa entre pasos o vaciada en cada paso. Junto con /channels/last-value conforman el trío del sistema de canales de LangGraph, y todas las interfaces provienen de /channels/base-channel. Un reducer como add_messages se declara en /graph/state-graph como Annotated[..., reducer], y en runtime entra a update mediante apply_writes de /pregel/algo. Como canal __pregel_tasks, Topic(accumulate=False) es el portador del fan-out Send, ver /subgraph/send-command.
Véase la documentación oficial: LangGraph docs · README