Skip to content

Capa de algoritmo de Pregel: interrupción, volcado de writes, programación

源码版本1.2.9

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_tasks son 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 en libs/langgraph/tests/unit/.
  • Compartidas entre síncrono y asíncrono: PregelLoop tiene dos subclases SyncPregelLoop y AsyncPregelLoop, 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_writes devuelve el conjunto updated_channels y prepare_next_tasks lo 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_interrupt deduplica con versions_seen: el versions_seen del canal INTERRUPT(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: cuando bump_step and updated_channels.isdisjoint(trigger_to_nodes), se llama channels[chan].finish()(finish:336-342) para enviar a todos los canales la señal «este es el último superpaso». Los canales permanentes (como LastValue) aprovechan para marcarse como no disponibles, de modo que prepare_next_tasks no dispare ningún nodo — es la «muerte natural» del grafo.

Archivos clave

  • should_interrupt:155-185 — comprueba si el grafo debe interrumpir; primero mira versions_seen[INTERRUPT] para ver si hubo novedades y luego si la tarea disparada cae en la lista interrupt_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-184interrupt_nodes == "*" significa que todas las tareas no ocultas cuentan; en caso contrario se filtra por task.name in interrupt_nodes.
  • apply_writes:232-345 — fusiona los writes de las tareas del paso en los canales y devuelve el conjunto updated_channels.
  • bump_step:256-259 — si alguna tarea tiene triggers, bump_step=True, lo que significa que el paso consume canales y avanza el contador step.
  • update seen versions:262-269 — cada tarea registra en versions_seen la 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 — llama channels[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-323channels[chan].update(vals) fusiona de verdad los writes en el canal; solo los canales is_available() entran en updated_channels.
  • bump_step notify:326-333 — con bump_step=True, los canales disponibles no actualizados en este paso también reciben update(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ía finish() 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 de Send) como PULL (de la activación por bordes) salen de aquí.
  • updated_channels optimization:475-486 — con el updated_channels del 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 por prepare_push_task_functional / prepare_push_task_send, PULL pasa por _triggers + _proc_input para 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; admite fresh=True para 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»:

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

python
# 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:

  1. 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, _triggers comparará con versiones futuras y solo se disparará si hubo cambio.
  2. consume() sobre los canales leídos — los canales tipo Topic, de «leer una vez y se limpia», vacían su búfer interno en consume() y solo avanza el número de versión. Esa es la razón por la que un mensaje fan-out con Send, una vez consumido, no vuelve a disparar a otros nodos.
  3. 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 canales is_available() (con valor y sin finish) entran en updated_channels — los canales tras finish() 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:

python
# 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 Send del canal TASKS; provienen del return Send(...) de algún nodo en el paso previo, y cada Send corresponde a una tarea PUSH.
  • Tarea PULL: hace lookup inverso sobre updated_channels con trigger_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. El sorted es 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_interrupt solo devuelve una lista no vacía cuando any_updates_since_prev_interrupt es verdadero, y la actualización de versions_seen[INTERRUPT] ocurre en la lógica implícita tras apply_writes dentro 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_step es la única señal de avance de step: si any(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) — mientras updated_channels tenga 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_channel contiene writes a un canal que no está en channels, solo se hace logger.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=True de local_read copia el canal: en lectura fresh se hace channels[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 pero apply_writes aún no ha corrido» sin afectar al estado global que ven otras tareas.
  • Algoritmo de hash del task id en prepare_single_task: cuando checkpoint["v"] > 1 se 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