Skip to content

PregelRunner: el ejecutor de tareas dentro de un superpaso

源码版本1.2.9

Responsabilidades

PregelRunner hace una sola cosa: tomar el lote de PregelExecutableTask que preparó PregelLoop.tick, ejecutarlo de verdad, recolectar sus writes y devolverlos a PregelLoop. No le importa el aspecto del grafo, ni el checkpoint, ni cómo reducen los canales — solo se ocupa de «ejecución concurrente + manejo de fallos + persistencia de writes».

Su posición es muy estrecha: dentro del bucle principal de Pregel.astream(sync main loop:2964-2984), cuando while loop.tick() devuelve True, justo después viene for _ in runner.tick([t for t in loop.tasks.values() if not t.writes], ...): y, al terminar, loop.after_tick(). Es decir, la entrada del runner es «tareas calculadas por tick que aún no han escrito writes» y la salida es «los writes de esas tareas ya están en checkpoint_pending_writes»; no hay un segundo búfer intermedio.

El propio PregelRunner es muy fino: dos métodos públicos tick(tick:176) y commit(commit:574), más una versión asíncrona atick. El trabajo pesado está en los invocadores de bajo nivel _call / _acall(_call:700)/(_acall:789) — son los responsables de envolver la función del nodo en un future programable, gestionar la política de retry, el fan-out de subtask Send, las llamadas de streaming chunk y, finalmente, volcar los writes al objeto task.

Motivación de diseño

¿Por qué el runner es un generator y no un método normal?

  • Yield en streaming: el runner tiene forma de generator for _ in runner.tick(...)(yield return type:188); cada vez que termina una tarea hace yield y devuelve el control al llamador — que así puede emitir de inmediato el streaming chunk al cliente, sin esperar a que termine el superpaso entero. Esa es la base del streaming de LangGraph.
  • Camino rápido de tarea única: con len(tasks) == 1 and timeout is None and get_waiter is None se hace una llamada síncrona directa, sin levantar thread pool(fast path:203-254). La mayoría de los escenarios agent de rama única caen aquí y se ahorran el overhead del thread pool.
  • commit como callback de Future: FuturesDict.on_done(on_done:116) se dispara al completar cada future e internamente es weakref.WeakMethod(self.commit)commit es el gancho de persistencia de writes por tarea, no el cierre batch. Así el orden con que los writes entran en pending_writes es naturalmente el de finalización, y el checkpoint puede escribirse cuanto antes.
  • Encaminamiento al error handler: _should_route_to_error_handler(_should_route_to_error_handler:171-174) comprueba si el nodo fallido tiene configurado un nodo error_handler. Si acierta, añade el id de la excepción a _handled_exception_ids y programa una tarea handler en sustitución de la tarea original — así el fallo no hace panic inmediato, sino que se cierra por otro nodo.
  • Cancelación con _should_stop_others: cuando cualquier tarea lanza una excepción que no es GraphBubbleUp, el runner cancela las demás tareas in-flight del mismo lote(_should_stop_others:616-634) — garantiza la semántica dentro del superpaso de «o todas, o al handler». GraphInterrupt se excluye explícitamente, porque una interrupción no cuenta como fallo.

Archivos clave

  • class PregelRunner:135-138 — definición de la clase; el docstring resume su responsabilidad en una frase: ejecutar tareas, commitar writes, ceder control y, si hace falta, interrumpir otras tareas.
  • __init__:140-169 — mantiene weakrefs de submit / put_writes, node_error_handler_map, callbacks schedule_error_handler/aschedule_error_handler y _handled_exception_ids, el conjunto de dedup de excepciones que sobrevive entre ticks.
  • FuturesDict:75-134 — subclase de dict a medida; el callback on_done dispara commit al completar cada future y el campo should_stop retiene un partial de _should_stop_others.
  • tick signature:176-188 — la entrada es Iterable[PregelExecutableTask] y la salida Iterator[None]; la signatura deja explícita la semántica de generator.
  • fast path:203-254 — con tarea única, sin timeout y sin waiter, se ejecuta síncrono directo; en fallo se programa el error handler si procede.
  • schedule tasks:259-276 — ruta multi-tarea: se lanza un future por tarea vía self.submit() (en realidad PregelLoop.submit).
  • concurrent wait loop:282-323 — bucle con concurrent.futures.wait(FIRST_COMPLETED); cada vez que termina una tarea se hace commit, se emite la salida y, si hace falta, se levanta una tarea handler.
  • commit method:574-613 — cierre por tarea: cancelled escribe ERROR, GraphInterrupt escribe INTERRUPT, excepción común escribe ERROR + ERROR_SOURCE_NODE, y una finalización normal escribe task.writes + el marcador NO_WRITES.
  • _should_stop_others:616-634 — decide si cancelar otras tareas; excluye explícitamente GraphBubbleUp y las excepciones ya handled.
  • _panic_or_proceed:650-697 — función de cierre: cancela todos los futures in-flight, fusiona varios GraphInterrupt en uno y, en timeout, lanza TimeoutError.
  • _call:700-787 — envoltorio de la llamada al nodo: retry, emisión de stream chunk, callback schedule_task para el fan-out de Send, acumulación de task.writes.
  • _should_route_to_error_handler:171-174 — devuelve True cuando task.name in self.node_error_handler_map; el propio nodo handler no puede ser encaminado (evita recursión).

Flujo de datos

Al inicio de tick se construye FuturesDict, un dict mejorado que llama automáticamente a commit al completar cada future:

python
def tick(
    self,
    tasks: Iterable[PregelExecutableTask],
    *,
    reraise: bool = True,
    timeout: float | None = None,
    retry_policy: Sequence[RetryPolicy] | None = None,
    get_waiter: Callable[[], concurrent.futures.Future[None]] | None = None,
    schedule_task: Callable[
        [PregelExecutableTask, int, Call | None],
        PregelExecutableTask | None,
    ],
) -> Iterator[None]:
    tasks = tuple(tasks)
    futures = FuturesDict(
        callback=weakref.WeakMethod(self.commit),
        event=threading.Event(),
        should_stop=partial(
            _should_stop_others, handled_exception_ids=self._handled_exception_ids
        ),
        future_type=concurrent.futures.Future,
    )
    # give control back to the caller
    yield

Esto viene de tick 头部:176-199. El primer yield devuelve el control de inmediato — la primera iteración de for _ in runner.tick(...) del llamador solo «arranca»; la ejecución real de tareas ocurre en iteraciones posteriores. Es una forma simplificada de corrutina con generator, para no introducir la palabra clave async en la ruta síncrona.

commit es la función de cierre por tarea, que vuelca los writes del task a pending_writes:

python
def commit(
    self,
    task: PregelExecutableTask,
    exception: BaseException | None,
) -> None:
    if isinstance(exception, asyncio.CancelledError):
        task.writes.append((ERROR, exception))
        self.put_writes()(task.id, task.writes)
    elif exception:
        if isinstance(exception, GraphInterrupt):
            if exception.args[0]:
                writes = [(INTERRUPT, exception.args[0])]
                if resumes := [w for w in task.writes if w[0] == RESUME]:
                    writes.extend(resumes)
                self.put_writes()(task.id, writes)
        elif isinstance(exception, GraphBubbleUp):
            pass
        else:
            task.writes.append((ERROR, exception))
            if self._should_route_to_error_handler(task) and not isinstance(
                exception, GraphBubbleUp
            ):
                task.writes.append((ERROR_SOURCE_NODE, task.name))
                self._handled_exception_ids.add(id(exception))
            self.put_writes()(task.id, task.writes)
    else:
        if self.node_finished and (
            task.config is None or TAG_HIDDEN not in task.config.get("tags", [])
        ):
            self.node_finished(task.name)
        if not task.writes:
            task.writes.append((NO_WRITES, None))
        self.put_writes()(task.id, task.writes)

Esto viene de commit:574-613. put_writes es una weakref a PregelLoop.put_writes(put_writes:415) que mete los writes en checkpoint_pending_writes y lanza el future del checkpointer para persistir. Hay que observar varios writes especiales:

  • CancelledError → escribe (ERROR, exception), de modo que al cerrar after_tick la tarea cuenta como «fallida pero registrada», sin panic.
  • GraphInterrupt → escribe (INTERRUPT, ...) + cualquier write RESUME; el valor de reanudación de la interrupción debe persistirse junto con ella y leerse junto al volver a resume.
  • Excepción común + error handler configurado → añade un write extra (ERROR_SOURCE_NODE, task.name); esta marca la barrerá PregelLoop._resume_error_handlers_if_applicable(ERROR_SOURCE_NODE scan:775) para programar el nodo handler correspondiente al reanudar.
  • Finalización normal sin writes → se añade (NO_WRITES, None) para que prepare_next_tasks no la tome como «aún no corrida»(NO_WRITES marker:609-611).

El flujo de ejecución completo del runner se puede dibujar así:

Límites y fallos

  • _handled_exception_ids persiste entre ticks: en __init__ este conjunto es de instancia, no por tick(_handled_exception_ids:169). Una vez que el id de un objeto excepción entra al conjunto, _should_stop_others no lo vuelve a tratar como fallo, evitando que tras correr el error handler se reprotege la excepción original.
  • GraphBubbleUp no cuenta como fallo: commit hace pass explícito sobre GraphBubbleUp, sin escribir nada(GraphBubbleUp pass:592-594). Estas excepciones (como ParentCommand) son señales al grafo padre, no errores reales.
  • NO_WRITES evita que la tarea se repita: si una tarea termina con normalidad pero no escribió nada, también se añade (NO_WRITES, None)(NO_WRITES:609-611). Si no, tras after_tick, prepare_next_tasks vería la tarea sin writes, la creería no corrida y la volvería a ejecutar al reanudar.
  • schedule_error_handler puede devolver None: el propio nodo handler puede no estar configurado o no poder construirse, en cuyo caso la tarea sigue la vía de pánico común(handler optional:230-248). El write ERROR_SOURCE_NODE de commit coexiste con esta ruta; _resume_error_handlers_if_applicable lo intentará de nuevo al reanudar.
  • El timeout no lanza GraphInterrupt: _panic_or_proceed hace if inflight: raise timeout_exc_cls("Timed out")(timeout:691-697) — el timeout es una excepción real que burbujea fuera de Pregel, no se traga como GraphInterrupt.
  • Recorte de traceback: antes del reraise se saltan los frames listados en EXCLUDED_FRAME_FNAMES(tb trim:241-247), filtrando los stack frames internos de langgraph para que el traceback que ve el usuario apunte directamente a su código de nodo.

Resumen

PregelRunner es una cáscara fina: el generator tick programa la ejecución concurrente y el commit por tarea vuelca los writes a pending_writes, con encaminamiento al error handler y cancelación entre medias. Entenderlo es entender «cómo se ejecutan de verdad las tareas dentro de un superpaso». La capa de bucle PregelLoop se ve en /pregel/loop; las tres funciones de algoritmo prepare_next_tasks / apply_writes / should_interrupt en /pregel/algo; el ensamblaje integral en /pregel/pregel; cómo se procesan las salidas en streaming entre los yield del runner se ve en /stream/run-stream.

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