Ejecución concurrente: orquestación de futures con @task
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_optionstoma del config el callbackCONFIG_KEY_CALL(el_call/_acallinyectado por el runner), envuelve la task en una estructuraCally, víaschedule_task, la mete en el canalTASKS(schedule PUSH task:713); la siguiente ronda deprepare_next_tasksla descubrirá de forma natural. - Abstracción unificada síncrona/ asíncrona:
SyncAsyncFuture(SyncAsyncFuture:253) es a la vez unconcurrent.futures.Futurey un objeto con__await__—su__await__sólo haceyielduna 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,
_callbusca primero enfuturessi ya existe un future con el mismo id y lo reutiliza; si la task ya terminó (next_task.writesno vacío), toma el valorRETURNdirectamente 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,AsyncBackgroundExecutorleeconfig["max_concurrency"]y levanta unasyncio.Semaphore; todas las coros de task se envuelven congated(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; mantienesubmit/put_writes/node_error_handler_map/schedule_error_handler.PregelRunner.tick:176— ejecuta síncronamente un lote de tareas; un bucle conconcurrent.futures.wait(FIRST_COMPLETED)va recogiendo futures done y commiteando writes.PregelRunner.atick:360— versión asíncrona;asyncio.waitsustituye alwaitsíncrono, con la misma semántica de iterador._call:700— callback síncrono paraCONFIG_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íaschedule_task(en realidadloop.accept_push/aaccept_push), mete la estructuraCallen el canalTASKS.SyncAsyncFuture:253— puente que es a la vez Future síncrono y awaitable asíncrono; su__await__hace un únicoyieldpara ceder el control._call_with_options:276— implementación de_TaskFunction.__call__; tomaCONFIG_KEY_CALLdel config y lo invoca.BackgroundExecutor:40— gestor de contexto de thread pool síncrono;submitusacopy_context()para que el hilo hijo herede los contextvar.AsyncBackgroundExecutor:122— gestor de contexto asíncrono, basado enasyncio.get_running_loop(), con semáforo opcionalmax_concurrency.gated:214—async def gated(semaphore, coro): primeroacquire()y luegorelease(), 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»:
# 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_taskLas 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 puedeawait, 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íachainFuturey sólo se hace visible cuando la tarea padre haceawait(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.submitusanext_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_concurrencyno afecta a la concurrencia entre nodos: el semáforo deAsyncBackgroundExecutorsólo acota los coros de task (gated), no cuántos nodos schedulea a la vezPregelRunner.tick—la concurrencia entre nodos la decideprepare_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,
_callcae en la rama «dedup already ran» y devuelve el future existente o el resultado tomado dewrites(dedup already ran:724), garantizando idempotencia. - Heredar contextvar:
BackgroundExecutor.submitusacopy_context()(copy_context:62) para que el hilo hijo herede los contextvar del hilo padre (comoRunnableConfig); 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