Skip to content

Ejecución concurrente: orquestación de futures con @task

源码版本1.2.9

Responsabilidades

La semántica de @task en runtime la asume PregelRunner —la llamada a una task no se ejecuta al instante, sino que entrega la función y los argumentos al runner actual, que decide cuándo y en qué hilo/corutina correrla. Esta capa se encarga de tres cosas: convertir la llamada en una tarea PUSH y meterla en la cola, ejecutar concurrentemente múltiples tareas disparadas dentro del mismo superpaso (superstep) y usar un future para devolver el resultado al llamador. PregelRunner (PregelRunner:135) es el centro de esta orquestación.

Se sitúa por debajo de PregelLoop y por encima de BackgroundExecutor / AsyncBackgroundExecutor. loop.tick calcula qué tareas debe correr el paso y se las pasa a runner.tick / atick (PregelRunner.tick:176); el runner hace submit de cada tarea al executor y, mediante un FuturesDict, sigue todos los futures in-flight, escribe sus writes cuando se completan y enruta las excepciones al nodo de tratamiento de errores.

Motivación de diseño

  • Una llamada a task es una tarea PUSH: dentro del cuerpo de un entrypoint, escribir task_a(); task_b() no exige declarar explícitamente un edge —_call_with_options toma del config el callback CONFIG_KEY_CALL (el _call / _acall inyectado por el runner), envuelve la task en una estructura Call y, vía schedule_task, la mete en el canal TASKS (schedule PUSH task:713); la siguiente ronda de prepare_next_tasks la descubrirá de forma natural.
  • Abstracción unificada síncrona/ asíncrona: SyncAsyncFuture (SyncAsyncFuture:253) es a la vez un concurrent.futures.Future y un objeto con __await__ —su __await__ sólo hace yield una vez y cede el control al runner exterior. En un entrypoint síncrono, task().result() bloquea para obtener el valor; en uno asíncrono, await task() devuelve la corutina —ambos comparten el mismo objeto future.
  • Planificación next_tick: cuando se programa una subtask, no se ejecuta enseguida: se usa __next_tick__=True (next_tick:768) para decirle al executor «corre en el siguiente tick» —así, las writes del tick actual se commitean primero y emiten sus eventos de stream antes de que arranquen las subtareas, garantizando que el orden del stream respeta el orden de llamada escrito en el código.
  • Deduplicación reentrante: cuando una misma task se llama varias veces, _call busca primero en futures si ya existe un future con el mismo id y lo reutiliza; si la task ya terminó (next_task.writes no vacío), toma el valor RETURN directamente de writes y construye un future ya resuelto (dedup already ran:724) —evita repetir subtareas al reintentar la tarea padre.
  • Semáforo max_concurrency: en la ruta async, AsyncBackgroundExecutor lee config["max_concurrency"] y levanta un asyncio.Semaphore; todas las coros de task se envuelven con gated(semaphore, coro) (gated:214), limitando cuántas tasks pueden estar in-flight simultáneamente dentro de un grafo y evitando que las llamadas a LLM colapsen el event loop.

Archivos clave

  • PregelRunner:135 — el ejecutor concurrente; mantiene submit / put_writes / node_error_handler_map / schedule_error_handler.
  • PregelRunner.tick:176 — ejecuta síncronamente un lote de tareas; un bucle con concurrent.futures.wait(FIRST_COMPLETED) va recogiendo futures done y commiteando writes.
  • PregelRunner.atick:360 — versión asíncrona; asyncio.wait sustituye al wait síncrono, con la misma semántica de iterador.
  • _call:700 — callback síncrono para CONFIG_KEY_CALL; convierte la llamada en tarea PUSH, deduplica y encadena futures.
  • _acall:789 — callback asíncrono, misma semántica, con manejo de coroutines.
  • schedule PUSH task:713 — vía schedule_task (en realidad loop.accept_push / aaccept_push), mete la estructura Call en el canal TASKS.
  • SyncAsyncFuture:253 — puente que es a la vez Future síncrono y awaitable asíncrono; su __await__ hace un único yield para ceder el control.
  • _call_with_options:276 — implementación de _TaskFunction.__call__; toma CONFIG_KEY_CALL del config y lo invoca.
  • BackgroundExecutor:40 — gestor de contexto de thread pool síncrono; submit usa copy_context() para que el hilo hijo herede los contextvar.
  • AsyncBackgroundExecutor:122 — gestor de contexto asíncrono, basado en asyncio.get_running_loop(), con semáforo opcional max_concurrency.
  • gated:214async def gated(semaphore, coro): primero acquire() y luego release(), envolviendo el coro de la task con un límite de concurrencia.

Flujo de datos

_call es el núcleo de la llamada a una task: traduce «dentro del cuerpo del entrypoint, llamar task_fn(arg)» en «programar una tarea PUSH en el runner»:

python
# schedule the next task, if the callback returns one
if next_task := schedule_task(
    task(),
    scratchpad.call_counter(),
    Call(
        func,
        input,
        retry_policy=retry_policy,
        cache_policy=cache_policy,
        callbacks=callbacks,
        timeout=timeout,
    ),
):
    if fut := next(
        (
            f
            for f, t in list(futures().items())
            if t is not None and t == next_task.id
        ),
        None,
    ):
        # if the parent task was retried,
        # the next task might already be running
        pass
    elif next_task.writes:
        # if it already ran, return the result
        fut = concurrent.futures.Future()
        ret = next((v for c, v in next_task.writes if c == RETURN), MISSING)
        # ...
    else:
        # schedule the next task
        fut = submit()(
            run_with_retry,
            next_task,
            retry_policy,
            configurable={
                CONFIG_KEY_CALL: partial(
                    _call,
                    weakref.ref(next_task),
                    futures=futures,
                    # ...
                ),
            },
            __reraise_on_exit__=False,
            __next_tick__=True,
        )
        SKIP_RERAISE_SET.add(fut)
        futures()[fut] = next_task

Las tres ramas deciden en orden: (1) la task con el mismo id ya está corriendo —reutiliza el future; (2) la task ya corrió (writes no vacío) —construye un future resuelto a partir del valor RETURN en writes; (3) task totalmente nueva —submit() entrega run_with_retry al BackgroundExecutor o AsyncBackgroundExecutor, e inyecta como CONFIG_KEY_CALL un _call recursivo en el config de la subtarea, de modo que las llamadas a task() dentro de la subtarea pasan por la misma maquinaria de scheduling.

El corazón de PregelRunner.tick es el bucle sobre FuturesDict + concurrent.futures.wait(FIRST_COMPLETED) (wait loop:281). Cada vez que un future se completa, se extrae la task correspondiente y se llama a commit para escribir sus writes en el loop; si la task tenía nodo de tratamiento de errores configurado, se hace submit de una tarea error handler adicional. El bucle sólo termina cuando todos los futures in-flight están completos (o queda sólo el future placeholder get_waiter) y devuelve el control a loop.after_tick, que vuelca los valores a los canales.

Límites y fallos

  • Llamar a una task async desde un contexto síncrono lanza error: al inicio de _call, if inspect.iscoroutinefunction(func): raise RuntimeError (sync async mismatch:716) —en un entrypoint síncrono no se puede await, así que las tasks asíncronas no pueden correr.
  • Las excepciones de subtasks burbujean al padre: SKIP_RERAISE_SET.add(fut) marca el future de la subtask como «no relanzar al terminar el tick»; la excepción se encadena al future padre vía chainFuture y sólo se hace visible cuando la tarea padre hace await (skip reraise:781). Así, la propagación de errores sigue la jerarquía de llamadas en lugar de ser lanzada directamente por el runner y romper el tick entero.
  • __next_tick__ en modo síncrono es ejecución diferida: BackgroundExecutor.submit usa next_tick(ctx.run, fn, ...) (next_tick sync:66) para encolar la tarea en la siguiente ronda de tick; en el executor asíncrono, __next_tick__ es noop, porque las tareas de asyncio ya se programan de forma asíncrona nativa.
  • max_concurrency no afecta a la concurrencia entre nodos: el semáforo de AsyncBackgroundExecutor sólo acota los coros de task (gated), no cuántos nodos schedulea a la vez PregelRunner.tick —la concurrencia entre nodos la decide prepare_next_tasks, el semáforo sólo hace throttling (semaphore:153).
  • Al reintentar la tarea padre, las subtareas no se re-ejecutan: si la tarea padre entra en retry, _call cae en la rama «dedup already ran» y devuelve el future existente o el resultado tomado de writes (dedup already ran:724), garantizando idempotencia.
  • Heredar contextvar: BackgroundExecutor.submit usa copy_context() (copy_context:62) para que el hilo hijo herede los contextvar del hilo padre (como RunnableConfig); sin eso, la subtarea no tendría acceso al config.

Resumen

La concurrencia de @task proviene de la orquestación de futures de PregelRunner: la abstracción central es SyncAsyncFuture, que tiende un puente entre síncrono y asíncrono; _call / _acall convierten la llamada a task en tarea PUSH, y BackgroundExecutor / AsyncBackgroundExecutor aportan los hilos/corutinas de ejecución. Para el detalle de cómo se schedulea una tarea PUSH, véase /subgraph/send-command; para la semántica de task dentro de un entrypoint, véase /func/entrypoint-task. Véase la documentación oficial: documentación de LangGraph · README