Skip to content

PregelLoop: el driver del bucle de superpasos

源码版本1.2.9

Responsabilidades

PregelLoop es el «metrónomo» del runtime de LangGraph. Pregel mismo solo se ocupa de ensamblar recursos y exponer las entradas invoke / astream; lo que de verdad corre el grafo vuelta a vuelta es una instancia de PregelLoop: mantiene el punto de control (checkpoint) actual, el conjunto de canales (channel), pending_writes, el contador step, la máquina de estados status y un diccionario tasks — todo el contexto de un único superpaso (superstep). El patrón exterior while loop.tick(): ... loop.after_tick() se sostiene sobre él.

Sus dos métodos centrales son tick(tick:599) y after_tick(after_tick:683). El primero cubre «antes de correr»: comprueba si se excede el límite de pasos, llama a prepare_next_tasks para cargar en self.tasks los nodos que ejecutará este paso, verifica interrupciones, atiende drain_requested y reintegra los pending_writes dejados por una reanudación previa sobre los nodos exitosos. El segundo cubre «después de correr»: recoge todos los task.writes, llama a apply_writes para fusionarlos en los canales, limpia los pending, persiste el checkpoint y luego decide si aplica interrupt_after. Ambos, uno antes y otro después, dejan al ciclo BSP «leer entrada → programar → ejecutar → volcar writes → persistir» encajado entre medias.

Dicho de otro modo, PregelLoop no ejecuta nodos por sí mismo — ese es el trabajo de PregelRunner; solo decide «a quién miro en este paso, cómo cierro tras correr y si sigo con el paso siguiente».

Motivación de diseño

¿Por qué partir la lógica del bucle en tick y after_tick en lugar de resolverlo con un único método step()?

  • Ceder la ejecución al llamador: cuando tick devuelve True, el llamador (el cuerpo de Pregel.astream, ver sync main loop:2964-2984) se encarga de impulsar runner.tick para que termine de correr ese lote de tareas y luego vuelve a llamar a loop.after_tick. Así PregelLoop no necesita retener una referencia a PregelRunner; la relación entre ambos es de colaboración acoplada de forma laxa, no de anidamiento.
  • Máquina de estados explícita: el campo status toma siete valores "input" / "pending" / "done" / "draining" / "interrupt_before" / "interrupt_after" / "out_of_steps"(status Literal:256-264). Cada estado corresponde a un motivo de salida; el código externo se apoya en él para decidir si resume, si raise o si break, sin tener que adivinar la semántica del valor de retorno.
  • Semántica clara de reanudación: is_replaying solo es verdadero en el primer tick — indica que este paso está «reproduciendo» tareas que ya habían terminado antes de la interrupción, y se apaga cuando after_tick termina(is_replaying = False:716). Junto con _reapply_writes_to_succeeded_nodes, que repega los writes exitosos rescatados del checkpoint sobre las tareas en memoria, se consigue que los nodos exitosos no se repitan y los fallidos pasen al error handler.
  • Drain es la vía de intervención externa: RunControl.drain_requested lo fija el consumidor del stream (por ejemplo, el cliente quiere parar a mitad); en cuanto tick lo ve pasa al estado draining y devuelve False, de modo que el while externo sale de forma natural, en lugar de matar el nodo a la fuerza a mitad de ejecución.
  • Escritura delta en modo exit y orden temporal del checkpoint: _delta_write_futs(_delta_write_futs:207) recolecta los futures de put_writes de todos los canales delta; _checkpointer_put_after_previous debe drenarlos antes de escribir el siguiente checkpoint — garantiza el orden causal de «los writes producen el checkpoint, no que el checkpoint tape los writes».

Archivos clave

  • class PregelLoop:158 — definición de la clase; los campos son casi todos attributes desnudos, sin envolver con property, para que las subclases SyncPregelLoop / AsyncPregelLoop los puedan sobreescribir directamente.
  • status Literal:256-264 — siete estados; "out_of_steps" se dispara al alcanzar el límite de pasos, "draining" por petición externa de parada, "done" es cierre normal.
  • tick method:599-681 — toda la lógica de un tick: comprobación de pasos, prepare_next_tasks, rama done para tareas vacías, drain, _reapply_writes_to_succeeded_nodes, should_interrupt.
  • out_of_steps:607-609 — cuando self.step > self.stop entra en estado out_of_steps y devuelve False; el while externo sale.
  • done branch:653-655 — cuando self.tasks está vacío entra en estado done; es la «muerte natural» del grafo.
  • draining branch:657-659 — cuando control.drain_requested es verdadero entra en estado draining.
  • reapply + resume handlers:662-664 — al reanudar, primero se reaplican los pending writes sobre los nodos exitosos y luego se llama a _resume_error_handlers_if_applicable para preparar tareas handler en los nodos fallidos.
  • interrupt_before:667-671 — llama a should_interrupt para comprobar si el nodo cae en la lista interrupt_before; si acierta, raise GraphInterrupt.
  • after_tick method:683-726 — recoger writes → apply_writes → emit values → limpiar pending → _put_checkpoint → comprobar interrupt_after.
  • _reapply_writes_to_succeeded_nodes:736-749 — reaplica los pending writes sobre las tareas en memoria, pero saltando las cuatro señales de control ERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME.
  • _resume_error_handlers_if_applicable:751-816 — recorre la marca ERROR_SOURCE_NODE y construye tareas handler para los nodos fallidos, de modo que el runner salte la tarea original y corra directamente el handler.
  • put_writes:415-508 — entrada de los writes que produce una tarea; dedup del canal de control, acumula null task y lanza el future put_writes del checkpointer.
  • SyncPregelLoop:1469 / AsyncPregelLoop:1722 — las dos subclases, que sobreescriben __enter__/__exit__, accept_push, schedule_error_handler y otras piezas que diferencian síncrono de asíncrono.

Flujo de datos

La entrada de tick hace primero la comprobación del límite de pasos y luego llama a prepare_next_tasks para calcular «qué tareas corren en el paso siguiente». Este paso es la semilla del superpaso entero — todas las acciones posteriores trabajan sobre ese lote de tareas.

python
def tick(self) -> bool:
    """Execute a single iteration of the Pregel loop.

    Returns:
        True if more iterations are needed.
    """

    # check if iteration limit is reached
    if self.step > self.stop:
        self.status = "out_of_steps"
        return False

    # prepare next tasks
    self.tasks = prepare_next_tasks(
        self.checkpoint,
        self.checkpoint_pending_writes,
        self.nodes,
        self.channels,
        self.managed,
        self.config,
        self.step,
        self.stop,
        for_execution=True,
        manager=self.manager,
        store=self.store,
        checkpointer=self.checkpointer,
        trigger_to_nodes=self.trigger_to_nodes,
        updated_channels=self.updated_channels,
        retry_policy=self.retry_policy,
        cache_policy=self.cache_policy,
    )

Esto viene de tick 头部:599-629. self.step > self.stop es la condición de salida por «límite de pasos» — stop se fija en PregelLoop.__init__ a recursion_limit; si la profundidad de recursión se excede, se para en out_of_steps. prepare_next_tasks toma el checkpoint actual, pending_writes, la tabla de nodos y la de canales y calcula qué tareas ejecutar en este paso, volcándolas en self.tasks. for_execution=True indica que estas tareas se van a correr de verdad (no es un dry-run para vista previa en streaming).

A continuación vienen tres condiciones de salida en paralelo — tareas vacías / drain / reanudar reaplicando writes — y luego should_interrupt:

python
# if no more tasks, we're done
if not self.tasks:
    self.status = "done"
    return False

if self.control is not None and self.control.drain_requested:
    self.status = "draining"
    return False

# if there are pending writes from a previous loop, apply them
if not self.is_replaying and self.checkpoint_pending_writes:
    self._reapply_writes_to_succeeded_nodes(self.tasks)
    self._resume_error_handlers_if_applicable()

# before execution, check if we should interrupt
if self.interrupt_before and should_interrupt(
    self.checkpoint, self.interrupt_before, self.tasks.values()
):
    self.status = "interrupt_before"
    raise GraphInterrupt()

Esto viene de tick 中段:652-671. La comprobación de interrupt_before se hace solo después de reaplicar los pending_writes en la reanudación — porque al reanudar la versión de los canales ya refleja los writes anteriores y should_interrupt puede decidir correctamente «si hubo nuevas actualizaciones desde el último interrupt».

after_tick se llama cuando el runner ha terminado de correr todas las tareas y se encarga del cierre:

python
def after_tick(self) -> None:
    # finish superstep
    writes = [w for t in self.tasks.values() for w in t.writes]
    self._delta_channels_with_overwrite.update(
        ch
        for ch, v in writes
        if isinstance(self.specs.get(ch), DeltaChannel) and _get_overwrite(v)[0]
    )
    # all tasks have finished
    self.updated_channels = apply_writes(
        self.checkpoint,
        self.channels,
        self.tasks.values(),
        self.checkpointer_get_next_version,
        self.trigger_to_nodes,
    )

Esto viene de after_tick 头部:683-698. Primero se aplanan los writes de todas las tareas y se anotan qué canales delta han sufrido overwrite (afecta al valor inicial del sparse replay); luego se llama a apply_writes para fusionar esos writes en los canales; el valor de retorno, updated_channels, lo usará prepare_next_tasks en el paso siguiente para acelerar «encontrar el siguiente nodo a disparar». Al final de after_tick también se vacía checkpoint_pending_writes, se pone is_replaying a False, se persiste el punto de control con _put_checkpoint({"source": "loop"}) y luego se ejecuta la comprobación de interrupt_after(after_tick 尾部:714-724).

El ciclo de vida completo de un superpaso se puede dibujar así:

Límites y fallos

  • out_of_steps no es un error: step > stop devuelve False directamente y status se fija en out_of_steps; el while externo sale de forma natural. Más arriba, Pregel decide según el status si lanzar RecursionError o cerrar en silencio — es semántica de «límite blando», no una vía de excepción(out_of_steps:607-609).
  • draining no puede emitir lifecycle events: _emit_graph_lifecycle_event rechaza explícitamente que se lo llame en estado draining(draining guard:384-385), porque drain es una parada iniciada por el usuario y no cuenta como interrupción / reanudación.
  • Al reanudar, los pending writes con señales de control deben saltarse: _reapply_writes_to_succeeded_nodes tiene que saltar ERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME(skip control signals:746-747); si no, las tareas que fallaron se tratarían como exitosas y los handler ya corridos se reactivarían.
  • Reprogramación del error handler al reanudar: _resume_error_handlers_if_applicable solo aplica a los nodos con error_handler_node configurado(handler_node check:791-793); los nodos fallidos sin handler se dejan tal cual y siguen la vía de pánico por _should_stop_others del runner.
  • Semántica de acumulación del null task en put_writes: los writes con NULL_TASK_ID no se sobreescriben, se acumulan(null task accumulate:422-431); sirve para input writes — writes «que no pertenecen a ningún nodo pero deben entrar en el checkpoint».
  • Los futures de los canales delta deben persistir antes que el siguiente checkpoint: _delta_write_futs recolecta los futures put_writes de todos los canales delta y se drena antes de persistir el siguiente checkpoint(_delta_write_futs:201-207) — invertir el orden haría que el sparse replay no viera los writes que produjeron el checkpoint.

Resumen

PregelLoop es un metrónomo de superpasos dirigido por máquina de estados que envuelve «preparar tareas → correr → cerrar» en torno a PregelRunner sin tocar la ejecución. Una vez entendidos los dos tramos tick / after_tick, los siete estados de status y el mecanismo de reanudación con is_replaying + _reapply_writes_to_succeeded_nodes, se tiene el eje principal del modelo de ejecución de LangGraph. Cómo hace el runner para correr de verdad las tareas se ve en /pregel/runner; los detalles internos de la capa de algoritmo prepare_next_tasks / apply_writes / should_interrupt en /pregel/algo; el ensamblaje integral en /pregel/pregel. La semántica de escritura en canales se ve en /channel/base-channel.

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