Skip to content

La abstracción BaseChannel: el canal como primitiva de estado

源码版本1.2.9

Responsabilidades

En el runtime de LangGraph, el «canal (channel)» es la unidad mínima del estado. La instancia Pregel que obtienes al compilar un StateGraph no mantiene directamente el dict del usuario, sino un conjunto de instancias BaseChannel —un canal por cada key—, y cada canal se encarga de almacenar su valor, fusionar escrituras, exponer lecturas y serializar su estado actual a forma de checkpoint. BaseChannel es la clase base abstracta común a todas las clases de canal (BaseChannel:19).

Ella misma no almacena datos (__slots__ solo declara key y typ, ver __slots__:22-26). Solo estipula cuatro miembros abstractos que las subclases deben implementar: dos propiedades de tipo ValueType / UpdateType y tres acciones update / get / from_checkpoint (abstract methods:60-99). Además hay cinco métodos con implementación por defecto —checkpoint / is_available / consume / finish / copy—, que o bien recurren a get() como fallback o devuelven False actuando como no-op.

Dicho de otro modo, BaseChannel fija de una sola vez como interfaz las tres cuestiones «cómo se almacena, cómo se muta y cómo se deja rastro del estado». Al motor Pregel le basta invocar update / get / checkpoint; no necesita saber si detrás hay un LastValue, un Topic o un BinaryOperatorAggregate haciendo la fusión. Ahí radica por qué el motor y la forma del estado quedan totalmente desacoplados.

Motivación de diseño

¿Por qué modelar el estado como objeto «canal» en vez de pasar un dict entre nodos?

  • Semántica unificada de fusión de escrituras: dentro del mismo superpaso (superstep), varios nodos pueden escribir concurrentemente en la misma key. Un dict no puede expresar reglas de fusión distintas como «lista se appenda, escalar se sobrescribe, contador se suma»; el canal recoge esa estrategia en un único update(values: Sequence[Update]) -> bool y la subclase decide cómo fold.
  • Vista inmutable dentro del paso: lo que lee un nodo es la instantánea al inicio del paso; los valores escritos por otros nodos en el mismo paso solo son visibles en el siguiente. get solo lee el valor actual; update se invoca exclusivamente en la frontera del paso, de forma unificada, mediante apply_writes (apply_writes:317-323), impidiendo por mecanismo cualquier condición de carrera.
  • Serializable: cada canal puede producir un valor serializable con checkpoint() (checkpoint:49-58) y reconstruirse con from_checkpoint (from_checkpoint:60-65). La familia BaseCheckpointSaver se apoya precisamente en estas dos interfaces para escribir el estado completo del hilo en SQLite/Postgres.
  • Hooks de ciclo de vida: consume y finish (consume / finish:101-121) dan a canales especiales (p. ej. Topic(accumulate=False) o LastValueAfterFinish) una oportunidad para mutarse en la frontera del paso o al final de la ejecución; el no-op por defecto indica que la mayoría de los canales no necesitan esta semántica.
  • Tipado autodescriptivo: las dos propiedades ValueType / UpdateType (ValueType / UpdateType:28-36) permiten que en tiempo de compilación y de ejecución se introspeccione «qué almacena y qué escrituras acepta este canal»; StateGraph las usa en _is_field_channel:1862-1887 para reconocer declaraciones tipo Annotated[..., SomeChannel].

Archivos clave

  • BaseChannel class:19-26 — definición de la clase base abstracta; parámetros genéricos Value / Update / Checkpoint; __slots__ solo declara key y typ.
  • ValueType / UpdateType:28-36 — dos property abstractas; las subclases las usan para declarar el tipo de valor almacenado y el tipo de valor de escritura.
  • checkpoint:49-58 — implementación por defecto: devuelve directamente self.get(); en canal vacío devuelve MISSING.
  • from_checkpoint:60-65 — método abstracto: construye un canal nuevo equivalente a partir de un valor de checkpoint; la subclase debe implementarlo.
  • get:69-73 — interfaz abstracta de lectura; un canal vacío lanza EmptyChannelError.
  • is_available:75-85 — implementación por defecto que invoca get() y captura EmptyChannelError; las subclases suelen sobreescribirla con una comprobación más barata.
  • update:89-99 — interfaz abstracta de escritura; Pregel la invoca una vez al final de cada superpaso, en orden arbitrario; devuelve True si el canal cambió.
  • consume / finish:101-121 — hooks de ciclo de vida; por defecto no-op y devuelven False.
  • apply_writes:315-345 — única entrada por la que Pregel recoge y fusiona todas las escrituras del paso a los canales; update / consume / finish se invocan aquí de forma unificada.
  • _is_field_channel:1862-1887 — al compilar, StateGraph usa isinstance(item, BaseChannel) para traducir Annotated[..., channel] a una instancia de canal.

Flujo de datos

Las escrituras al canal ocurren dentro de apply_writes: Pregel agrega los writes de todas las tareas del paso por canal y luego invoca update (apply_writes:317-323):

python
# Apply writes to channels
updated_channels: set[str] = set()
for chan, vals in pending_writes_by_channel.items():
    if chan in channels:
        if channels[chan].update(vals) and next_version is not None:
            checkpoint["channel_versions"][chan] = next_version
            # unavailable channels can't trigger tasks, so don't add them
            if channels[chan].is_available():
                updated_channels.add(chan)

# Channels that weren't updated in this step are notified of a new step
if bump_step:
    for chan in channels:
        if channels[chan].is_available() and chan not in updated_channels:
            if channels[chan].update(EMPTY_SEQ) and next_version is not None:
                checkpoint["channel_versions"][chan] = next_version
                if channels[chan].is_available():
                    updated_channels.add(chan)

Nótese el segundo bloque: los canales no escritos por ninguna tarea en este paso también reciben una llamada update(EMPTY_SEQ). Es un contrato implícito de la interfaz BaseChannel —el update de la subclase debe aceptar una secuencia vacía, normalmente devolviendo False como señal de «sin cambios». La semántica «vaciar en cada paso» de Topic(accumulate=False) se apoya precisamente en esta llamada vacía: el update vacío le permite descartar los valores del paso anterior y que dejen de ser visibles en el siguiente paso.

La posición del ciclo de vida completo del canal dentro de un tick de Pregel:

Límites y fallos

  • EmptyChannelError: lo lanza get() cuando el canal nunca ha sido escrito (get:70-73). La implementación por defecto de is_available se apoya en capturar esta excepción para decidir la disponibilidad; las subclases que puedan decidir más barato deberían sobreescribirla.
  • Centinela MISSING: cuando get() lanza EmptyChannelError, checkpoint() devuelve langgraph._internal._typing.MISSING en lugar de None (checkpoint fallback:55-58), de modo que None puede aparecer como valor legítimamente almacenado sin confundirse con «vacío».
  • update devuelve False: significa que el canal no cambió, así que Pregel no actualizará channel_versions y por tanto prepare_next_tasks no disparará ningún nodo por este canal (update call:319); es la base del criterio de convergencia de Pregel.
  • Llamada con secuencia vacía: apply_writes también invoca update(EMPTY_SEQ) sobre los canales no modificados (EMPTY_SEQ update:329). El contrato de BaseChannel.update es que una entrada vacía devuelva False, pero Topic(accumulate=False) en este caso se vacía y puede devolver True, así que las subclases no pueden asumir que «entrada vacía = no hacer nada».
  • consume y finish por defecto son no-op (consume / finish:101-121): devolver False implica que la mayoría de los canales no participa en la gestión del ciclo de vida. apply_writes invoca finish al final del bump_step (finish call:338); solo canales especiales como LastValueAfterFinish lo usan para implementar la semántica «visible solo al final de la ejecución».
  • copy por defecto pasa por checkpoint: BaseChannel.copy por defecto hace self.from_checkpoint(self.checkpoint()) (copy:40-47). Si en una subclase el checkpoint es costoso (p. ej. precisa deep-copy de una lista grande), conviene sobreescribir copy con una implementación más barata; Topic.copy hace exactamente eso.

Resumen

BaseChannel fija de una vez como interfaz «cómo se almacena, cómo se fusiona y cómo se serializa el estado»; al motor Pregel le basta invocar en el orden fijo update → is_available → get → checkpoint → finish para correr sobre cualquier implementación de canal. Para entender las implementaciones concretas, consulta /channels/last-value y /channels/topic-binop; el proceso de lectura y escritura de canales por parte de Pregel se describe en /pregel/algo; cómo se persiste el valor devuelto por checkpoint() se explica en /checkpoint/base-saver.

Véase la documentación oficial: LangGraph docs · README