Skip to content

Topic y BinaryOperatorAggregate: canal acumulativo y reductor

源码版本1.2.9

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 — messages se appenda deduplicando por ID, counter se suma, tags se une por unión. En vez de fabricar una clase de canal para cada caso, mejor un canal genérico que acepte un operador binario operator: (Value, Value) -> Value y 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 de Annotated / Required / NotRequired y luego instancia typ() para obtener el valor por defecto. Tipos abstractos como collections.abc.Sequence no se pueden instanciar directamente, así que en typ fallback:83-88 se mapean a list / set / dict (typ instantiation:82-92).
  • Topic.accumulate en doble forma: con accumulate=True es una lista acumulativa entre pasos — update hace extend de los nuevos valores (update:77-85); con accumulate=False primero 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 canal TASKS: Send es válido dentro de un paso y caduca tras ser consumido.
  • _flatten admite escrituras mixtas: Topic.update acepta una secuencia mixta Value | list[Value] (_flatten:15-20), de modo que un nodo puede escribir un Send o [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 dataclass Overwrite (Overwrite:980) o el dict {"__overwrite__": value} (OVERWRITE key:95) son reconocidos por _get_overwrite (_get_overwrite:31-51) y, dentro de update, saltan el reducer y asignan directamente. Solo se permite un Overwrite por paso; un segundo lanza InvalidUpdateError (overwrite dedup:132-138).

Archivos clave

  • Topic class:23-41 — definición de la clase; parámetros genéricos Sequence[Value] / Value | list[Value] / list[Value]; el constructor acepta el flag accumulate.
  • _flatten:15-20 — función utilitaria: aplana una secuencia mixta Value | list[Value] a Iterator[Value].
  • Topic.copy:56-61 — sobreescribe copy para 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 directamente self.values; from_checkpoint sigue siendo compatible con el formato antiguo en forma de tuple.
  • Topic.update:77-85 — lógica central: con accumulate=False primero vacía y luego extiende; con accumulate=True extiende directamente; una secuencia flat vacía devuelve False.
  • Topic.get / is_available:87-94 — lista vacía lanza EmptyChannelError; is_available se sobreescribe como bool(self.values).
  • Binop class:65-92 — definición de la clase + __init__: acepta el callable binario operator e inicializa el valor con typ().
  • _get_overwrite:31-51 — reconoce tres representaciones de Overwrite (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 invocan operator(self.value, value) uno a uno; el bypass Overwrite es la excepción.
  • _is_field_binop:1890-1908 — en compilación reconoce el callable binario en los metadatos de Annotated[..., reducer] y lo instancia como BinaryOperatorAggregate.
  • TASKS channel:805-809Pregel.__init__ fija el canal __pregel_tasks como Topic(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):

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

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

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

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

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

Nó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 Overwrite coexisten: overwrite forms:44-50 acepta simultáneamente el dataclass Overwrite, 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 Overwrite por paso: overwrite dedup:132-138 lanza InvalidUpdateError ante una segunda ocurrencia, porque la semántica de «sobrescritura» es contradictoria si aparece varias veces en un mismo paso.
  • Los valores no-overwrite tras Overwrite se descartan en silencio: discard after overwrite:142-143; if not seen_overwrite decide si invocar el reducer; una vez seen_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-62 trata 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) devuelve True ante escritura vacía: clear returns True:79-80 — siempre que el valor viejo no estuviera vacío, vaciar cuenta como cambio. Esto hará que apply_writes actualice channel_versions y dispare los nodos que dependen de ese canal. Si usas Topic por tu cuenta, sé consciente: no maneja la escritura vacía en silencio como LastValue.
  • Topic compatible con checkpoints antiguos: tuple compatibility:70-74 comprueba isinstance(checkpoint, tuple); en versiones antiguas el checkpoint almacenaba un tuple (typ, values) y en la nueva versión se almacena directamente una list. Este bloque existe para no romper los archivos antiguos.
  • La firma del reducer se fuerza a dos parámetros: reducer signature check:1894-1907 usa inspect.signature para contar los parámetros posicionales; si no son exactamente dos, lanza error. Funciones con *args u otros variadic no se aceptan — fija «el reducer debe ser estrictamente (a, b) -> c».
  • add_messages es una factoría: _add_messages_wrapper:41-57 la decora con @_add_messages_wrapper; la función real puede invocarse directamente como add_messages(left, right) o como add_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