Capa de algoritmo de Pregel: interrupción, volcado de writes, programación
Responsabilidades
_algo.py es el «núcleo matemático» de PregelLoop y PregelRunner. El loop se ocupa de «cuántas vueltas dar», el runner de «cómo ejecutar concurrentemente», y algo se ocupa de «la semántica de cada paso»: cuándo interrumpir, cómo se fusionan los writes del nodo en los canales (channel), qué nodos se disparan en el paso siguiente. No tiene estado ni IO — todo son funciones puras; la entrada es checkpoint + canales + tareas, y la salida es una tabla nueva de versiones de canal / un conjunto updated_channels / un diccionario nuevo de tareas.
Sus tres funciones centrales se corresponden exactamente con las tres llamadas de PregelLoop.tick + after_tick: should_interrupt(should_interrupt:155) se llama en el tramo medio de tick y decide si lanzar GraphInterrupt antes de ejecutar; apply_writes(apply_writes:232) se llama al inicio de after_tick para fusionar los writes de todas las tareas del paso en los canales y devolver updated_channels; prepare_next_tasks(prepare_next_tasks:392) se llama al inicio de tick y, a partir del checkpoint actual, calcula qué tareas debe ejecutar este paso.
Dicho de otro modo, estas tres funciones determinan «qué ve, qué hace y en qué quedan los canales al cerrar cada paso dentro del ciclo BSP de Pregel». Entenderlas es entender las «reglas de cambio de estado» de Pregel.
Motivación de diseño
¿Por qué descomponer la lógica del algoritmo en funciones puras?
- Testabilidad:
should_interrupt/apply_writes/prepare_next_tasksson funciones puras con entrada (dict de checkpoint + mapping de canales) y salida (estructuras de datos). Se pueden testear unitariamente sin tocar el IO de Pregel, ni submit, ni checkpointer — es exactamente la forma de la mayoría de los tests enlibs/langgraph/tests/unit/. - Compartidas entre síncrono y asíncrono:
PregelLooptiene dos subclasesSyncPregelLoopyAsyncPregelLoop, pero la capa algo no tiene side effect y ambas subclases comparten la misma lógica(tick calls prepare_next_tasks:612). Si no, habría que escribir el algoritmo dos veces y parchear el bug en ambos lados. - Optimización con
updated_channels:apply_writesdevuelve el conjunto updated_channels yprepare_next_taskslo recibe como hint, saltándose «escanear todos los nodos para ver si el trigger dispara»(updated_channels hint:475-486). En grafos grandes, este paso reduce el escaneo de trigger O(N×M) a O(updated×triggered); es la optimización de rendimiento más importante. should_interruptdeduplica conversions_seen: elversions_seendel canalINTERRUPT(versions_seen INTERRUPT:163) registra la versión de canal vista en el último interrupt. Solo «si hubo nuevas actualizaciones de canal desde el último interrupt» se considera volver a interrumpir — evita el bucle infinito de disparar interrupt, luego resume y luego interrupt otra vez.- La señal finish al final de
apply_writes: cuandobump_step and updated_channels.isdisjoint(trigger_to_nodes), se llamachannels[chan].finish()(finish:336-342) para enviar a todos los canales la señal «este es el último superpaso». Los canales permanentes (comoLastValue) aprovechan para marcarse como no disponibles, de modo queprepare_next_tasksno dispare ningún nodo — es la «muerte natural» del grafo.
Archivos clave
should_interrupt:155-185— comprueba si el grafo debe interrumpir; primero miraversions_seen[INTERRUPT]para ver si hubo novedades y luego si la tarea disparada cae en la listainterrupt_nodes.any_updates_since_prev_interrupt:161-168— compara versiones de canal para decidir «si hubo actualizaciones desde el último interrupt»; núcleo de la dedup de interrupciones.interrupt_nodes filter:170-184—interrupt_nodes == "*"significa que todas las tareas no ocultas cuentan; en caso contrario se filtra portask.name in interrupt_nodes.apply_writes:232-345— fusiona los writes de las tareas del paso en los canales y devuelve el conjuntoupdated_channels.bump_step:256-259— si alguna tarea tienetriggers,bump_step=True, lo que significa que el paso consume canales y avanza el contador step.update seen versions:262-269— cada tarea registra enversions_seenla versión de los canales trigger que leyó; así, en la próxima evaluación de_triggers, lo que se compara es «si hubo nuevas actualizaciones desde la última ejecución».consume channels:284-292— llamachannels[chan].consume()para marcar como consumidos los canales leídos en este paso; es la implementación de la semántica de canal PULL «leer una vez y se limpia».apply writes to channels:315-323—channels[chan].update(vals)fusiona de verdad los writes en el canal; solo los canalesis_available()entran enupdated_channels.bump_step notify:326-333— conbump_step=True, los canales disponibles no actualizados en este paso también recibenupdate(EMPTY_SEQ)para avanzar su número de versión — así, los nodos suscritos a esos canales pero no disparados en este paso tampoco se dispararán por error en el siguiente tick.finish signal:336-342— si ningún nodo fue disparado por actualización en este paso, se envíafinish()a todos los canales; los permanentes aprovechan para volverse unavailable y el grafo se cierra solo.prepare_next_tasks:392-513— calcula las tareas a correr en el paso siguiente; tanto PUSH (del fan-out deSend) como PULL (de la activación por bordes) salen de aquí.updated_channels optimization:475-486— con elupdated_channelsdel paso previo +trigger_to_nodes, hace lookup inverso del conjunto de nodos disparados y se salta el escaneo completo.prepare_single_task:524— construcción de una tarea individual; PUSH pasa porprepare_push_task_functional/prepare_push_task_send, PULL pasa por_triggers+_proc_inputpara sacar la entrada._triggers:1260-1277— decide si el canal trigger de un nodo «tiene una versión nueva sin leer»; es la condición central de la programación PULL.local_read:188-224— entrada para que la función del nodo lea valores del canal; admitefresh=Truepara leer la vista local «acabo de escribir pero aún no está fusionada en el canal», usado por la semántica de escritura y lectura inmediata del borde condicional (conditional edge).
Flujo de datos
should_interrupt es la más sencilla de las tres; la lógica es «hay nuevas actualizaciones + la tarea actual cae en la lista»:
def should_interrupt(
checkpoint: Checkpoint,
interrupt_nodes: All | Sequence[str],
tasks: Iterable[PregelExecutableTask],
) -> list[PregelExecutableTask]:
"""Check if the graph should be interrupted based on current state."""
version_type = type(next(iter(checkpoint["channel_versions"].values()), None))
null_version = version_type() # type: ignore[misc]
seen = checkpoint["versions_seen"].get(INTERRUPT, {})
# interrupt if any channel has been updated since last interrupt
any_updates_since_prev_interrupt = any(
version > seen.get(chan, null_version) # type: ignore[operator]
for chan, version in checkpoint["channel_versions"].items()
)
# and any triggered node is in interrupt_nodes list
return (
[task for task in tasks if ...]
if any_updates_since_prev_interrupt
else []
)Esto viene de should_interrupt:155-185. La clave es versions_seen[INTERRUPT], un registro especial — no es la versión vista por un nodo, sino «el snapshot de las versiones de todos los canales en el momento del último interrupt». any_updates_since_prev_interrupt compara ese snapshot con el channel_versions actual; basta con que cualquier canal tenga versión más alta para considerar «hay algo nuevo» y recién entonces pasar al juicio de interrupción. Así, al reanudar, aunque el siguiente tick dispare nodos con el mismo nombre, no se vuelve a interrumpir inmediatamente — hay que esperar a que el canal tenga writes realmente nuevos.
apply_writes es el núcleo del cierre; primero actualiza versions_seen, luego consume los canales leídos y finalmente escribe los nuevos valores:
# update seen versions
for task in tasks:
checkpoint["versions_seen"].setdefault(task.name, {}).update(
{
chan: checkpoint["channel_versions"][chan]
for chan in task.triggers
if chan in checkpoint["channel_versions"]
}
)
# Consume all channels that were read
for chan in {
chan
for task in tasks
for chan in task.triggers
if chan not in RESERVED and chan in channels
}:
if channels[chan].consume() and next_version is not None:
checkpoint["channel_versions"][chan] = next_version
# 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
if channels[chan].is_available():
updated_channels.add(chan)Esto viene de apply_writes 核心:262-323. La semántica de los tres tramos es clara:
- Actualizar
versions_seen— anota la versión actual de los canales trigger de cada tarea como prueba de que «este nodo ya vio esta versión de canal». La próxima vez,_triggerscomparará con versiones futuras y solo se disparará si hubo cambio. consume()sobre los canales leídos — los canales tipoTopic, de «leer una vez y se limpia», vacían su búfer interno enconsume()y solo avanza el número de versión. Esa es la razón por la que un mensaje fan-out conSend, una vez consumido, no vuelve a disparar a otros nodos.update(vals)escribe nuevos valores — recolecta los writes de todas las tareas del paso para cada canal y los fusiona de una sola vez, disparando la lógica de merge de canales reducer /LastValue. Solo los canalesis_available()(con valor y sin finish) entran enupdated_channels— los canales trasfinish()avanzan de versión pero no disparan ningún nodo.
prepare_next_tasks se llama al inicio de tick y decide qué tareas corre el paso:
# Consume pending tasks
tasks_channel = cast(Topic[Send] | None, channels.get(TASKS))
if tasks_channel and tasks_channel.is_available():
for idx, _ in enumerate(tasks_channel.get()):
if task := prepare_single_task(
(PUSH, idx), None, ..., for_execution=for_execution, ...
):
tasks.append(task)
# This section is an optimization that allows which nodes will be active
# during the next step.
if updated_channels and trigger_to_nodes:
triggered_nodes: set[str] = set()
for channel in updated_channels:
if node_ids := trigger_to_nodes.get(channel):
triggered_nodes.update(node_ids)
candidate_nodes: Iterable[str] = sorted(triggered_nodes)
elif not checkpoint["channel_versions"]:
candidate_nodes = ()
else:
candidate_nodes = processes.keys()
for name in candidate_nodes:
if task := prepare_single_task((PULL, name), None, ..., for_execution=for_execution, ...):
tasks.append(task)
return {t.id: t for t in tasks}Esto viene de prepare_next_tasks 核心:441-513. Aquí confluyen dos tipos de tareas:
- Tarea PUSH: saca objetos
Senddel canalTASKS; provienen delreturn Send(...)de algún nodo en el paso previo, y cadaSendcorresponde a una tarea PUSH. - Tarea PULL: hace lookup inverso sobre
updated_channelscontrigger_to_nodes— ese mapping se calcula en tiempo de compilación como «canal → lista de nodos suscritos»; la unión directa da los nodos candidatos a disparar en este paso. Elsortedes para que el orden de tareas sea determinista; no afecta al resultado de la ejecución, pero sí a la repetición del checkpoint.
Cada nodo candidato pasa luego por _triggers para una segunda confirmación(_triggers call:606-612) — solo «el canal tiene una versión nueva y el nodo aún no la ha visto» genera de verdad una tarea. Este paso es la clave de la deduplicación de nodos PULL: aunque updated_channels incluya el canal trigger de un nodo, si versions_seen[name] ya registró esa versión, la tarea no se reprograma.
La posición de toda la capa algo dentro del ciclo Pregel:
Límites y fallos
- Si
versions_seen[INTERRUPT]no se actualiza, hay bucle infinito:should_interruptsolo devuelve una lista no vacía cuandoany_updates_since_prev_interruptes verdadero, y la actualización deversions_seen[INTERRUPT]ocurre en la lógica implícita trasapply_writesdentro de_put_checkpoint. Si al reanudar tras un interrupt no se avanza correctamente la versión vista del canal INTERRUPT, la misma actualización se considera «nueva» una y otra vez, provocando un bucle interrupt-resume-interrupt(versions_seen INTERRUPT:163-168). bump_stepes la única señal de avance de step: siany(t.triggers for t in tasks)es falso, no se hace bump — esta es la semántica de «null task solo escribe canales sin avanzar step»(bump_step:259). El null task se usa para inyectar input writes antes de tick y no debe mover el contador step.finish()solo se llama cuando «ningún nodo se disparó»:bump_step and updated_channels.isdisjoint(trigger_to_nodes)(finish condition:336) — mientrasupdated_channelstenga algún canal suscrito por un nodo, no se hace finish. Significa que finish se dispara solo «al cerrar el último paso del grafo», no en cada paso.- Escribir en un canal desconocido solo avisa, no falla: si
pending_writes_by_channelcontiene writes a un canal que no está enchannels, solo se hacelogger.warning(unknown channel warn:310-313). Es por tolerancia — por ejemplo, si un nodo escribe en un canal ya terminado, no debe tirar todo el grafo. - La ruta
fresh=Truedelocal_readcopia el canal: en lecturafreshse hacechannels[k].copy()por cada canal(fresh read copy:213-219) y se aplican sobre la copia los writes de esta tarea — así, los nodos de borde condicional (conditional edge) pueden leer la vista local «acabo de escribir peroapply_writesaún no ha corrido» sin afectar al estado global que ven otras tareas. - Algoritmo de hash del task id en
prepare_single_task: cuandocheckpoint["v"] > 1se usa_xxhash_str, si no_uuid5_str(task_id_func:550); es por compatibilidad entre versiones de checkpoint — los task ids generados con uuid5 en checkpoints antiguos deben poder reproducirse.
Resumen
_algo.py tres funciones puras should_interrupt / apply_writes / prepare_next_tasks extraen limpiamente la capa semántica de Pregel: cuándo interrumpir, cómo fusionar los writes, quién corre en el siguiente paso. No tienen estado ni IO y, llamadas desde PregelLoop.tick + after_tick, impulsan todo el ciclo BSP. La capa de impulso del bucle en /pregel/loop; la capa de ejecución en /pregel/runner; el ensamblaje integral en /pregel/pregel; la semántica de update / consume / finish / is_available de los canales en /channel/base-channel.
Véase la documentación oficial: documentación de LangGraph · README