La abstracción BaseChannel: el canal como primitiva de estado
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]) -> booly 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.
getsolo lee el valor actual;updatese invoca exclusivamente en la frontera del paso, de forma unificada, medianteapply_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 confrom_checkpoint(from_checkpoint:60-65). La familiaBaseCheckpointSaverse apoya precisamente en estas dos interfaces para escribir el estado completo del hilo en SQLite/Postgres. - Hooks de ciclo de vida:
consumeyfinish(consume / finish:101-121) dan a canales especiales (p. ej.Topic(accumulate=False)oLastValueAfterFinish) 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»;StateGraphlas usa en_is_field_channel:1862-1887para reconocer declaraciones tipoAnnotated[..., SomeChannel].
Archivos clave
BaseChannel class:19-26— definición de la clase base abstracta; parámetros genéricosValue / Update / Checkpoint;__slots__solo declarakeyytyp.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 directamenteself.get(); en canal vacío devuelveMISSING.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 lanzaEmptyChannelError.is_available:75-85— implementación por defecto que invocaget()y capturaEmptyChannelError; 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; devuelveTruesi el canal cambió.consume / finish:101-121— hooks de ciclo de vida; por defecto no-op y devuelvenFalse.apply_writes:315-345— única entrada por la que Pregel recoge y fusiona todas las escrituras del paso a los canales;update/consume/finishse invocan aquí de forma unificada._is_field_channel:1862-1887— al compilar,StateGraphusaisinstance(item, BaseChannel)para traducirAnnotated[..., 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):
# 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 lanzaget()cuando el canal nunca ha sido escrito (get:70-73). La implementación por defecto deis_availablese apoya en capturar esta excepción para decidir la disponibilidad; las subclases que puedan decidir más barato deberían sobreescribirla.- Centinela
MISSING: cuandoget()lanzaEmptyChannelError,checkpoint()devuelvelanggraph._internal._typing.MISSINGen lugar deNone(checkpoint fallback:55-58), de modo queNonepuede aparecer como valor legítimamente almacenado sin confundirse con «vacío». updatedevuelveFalse: significa que el canal no cambió, así que Pregel no actualizaráchannel_versionsy por tantoprepare_next_tasksno 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_writestambién invocaupdate(EMPTY_SEQ)sobre los canales no modificados (EMPTY_SEQ update:329). El contrato deBaseChannel.updatees que una entrada vacía devuelvaFalse, peroTopic(accumulate=False)en este caso se vacía y puede devolverTrue, así que las subclases no pueden asumir que «entrada vacía = no hacer nada». consumeyfinishpor defecto son no-op (consume / finish:101-121): devolverFalseimplica que la mayoría de los canales no participa en la gestión del ciclo de vida.apply_writesinvocafinishal final delbump_step(finish call:338); solo canales especiales comoLastValueAfterFinishlo usan para implementar la semántica «visible solo al final de la ejecución».copypor defecto pasa por checkpoint:BaseChannel.copypor defecto haceself.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 sobreescribircopycon una implementación más barata;Topic.copyhace 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