Exécution concurrente : orchestration de futures par @task
Responsabilités
La sémantique d'exécution de @task est prise en charge par PregelRunner — l'appel d'une task ne s'exécute pas immédiatement, il remet la fonction et les arguments au runner courant, qui décide quand et sur quel thread / coroutine l'exécuter. Cette couche est responsable de trois choses : convertir l'appel de task en une tâche PUSH insérée dans la file, exécuter en concurrence plusieurs tasks déclenchées au cours d'un même superpas (superstep), et utiliser un future pour renvoyer le résultat à l'appelant. PregelRunner(PregelRunner:135) est le centre de cette orchestration.
Sa position est sous PregelLoop, au-dessus de BackgroundExecutor / AsyncBackgroundExecutor. loop.tick calcule quelles tâches exécuter à cette étape et les remet à runner.tick / atick(PregelRunner.tick:176) ; le runner les soumet à l'executor via submit, suit tous les futures en vol avec FuturesDict, et à chaque future terminé écrit les writes en retour et route les exceptions vers le nœud de gestion d'erreurs.
Motivation de conception
- Appel de task = tâche PUSH : dans le corps de la fonction entrypoint, écrire
task_a(); task_b()ne nécessite pas de déclarer d'arête —_call_with_optionsrécupère la callbackCONFIG_KEY_CALLdepuis la config (c'est le_call/_acallinjecté par le runner), emballe la task en une structureCall, et l'insère dans le canalTASKSviaschedule_task(schedule PUSH task:713) ;prepare_next_tasksau tour suivant la trouvera naturellement. - Abstraction unifiée sync / async :
SyncAsyncFuture(SyncAsyncFuture:253) est à la fois unconcurrent.futures.Futureet un objet avec__await__— son__await__fait un seulyield, rendant le contrôle au runner supérieur. Dans une entrypoint synchrone,task().result()bloque pour récupérer la valeur ; dans une entrypoint asynchrone,await task()cède la coroutine ; les deux partagent le même objet future. - Ordonnancement next_tick : quand une sous-task est ordonnancée, elle ne s'exécute pas immédiatement ; on utilise
__next_tick__=True(next_tick:768) pour dire à l'executor « ne l'exécute qu'au prochain tick » — ainsi les writes du tick courant sont d'abord commités et les événements stream émis avant la sous-task, pour garantir que l'ordre du stream correspond à l'ordre d'appel écrit dans le code. - Déduplication par ré-entrée : quand la même task est appelée plusieurs fois,
_callcherche d'abord dansfuturesun future de même id, et le réutilise s'il existe ; si la task a déjà fini (next_task.writesnon vide), il récupère directement la valeurRETURNdepuis les writes et construit un future déjà complété(dedup already ran:724) — pour éviter qu'une sous-task soit ré-exécutée quand le père est retenté. - Sémaphore
max_concurrency: dans le chemin asynchrone,AsyncBackgroundExecutorlitconfig["max_concurrency"]et instancie unasyncio.Semaphore; toutes les coros de task sont emballées dansgated(semaphore, coro)(gated:214) pour plafonner le nombre de tasks en vol simultanées dans le graphe et empêcher les appels LLM de plomber l'event loop.
Fichiers clés
PregelRunner:135— exécuteur concurrent, détientsubmit/put_writes/node_error_handler_map/schedule_error_handler.PregelRunner.tick:176— exécute un lot de tâches en synchrone ; boucleconcurrent.futures.wait(FIRST_COMPLETED)pour récupérer les futures terminés et commit les writes.PregelRunner.atick:360— version asynchrone,asyncio.waitremplace lewaitsynchrone, sémantique d'itérateur identique._call:700— callbackCONFIG_KEY_CALLsynchrone, responsable de convertir l'appel de task en tâche PUSH, dédupliquer, et chaîner le future._acall:789— callback asynchrone, même sémantique, traite les coroutines.schedule PUSH task:713— viaschedule_task(loop.accept_push/aaccept_push) insère la structureCalldans le canalTASKS.SyncAsyncFuture:253— objet pont à la fois future synchrone et awaitable asynchrone ;__await__fait un seulyieldpour rendre le contrôle._call_with_options:276— implémentation de_TaskFunction.__call__, récupèreCONFIG_KEY_CALLdepuis la config et l'invoque.BackgroundExecutor:40— gestionnaire de pool de threads synchrones,submitutilisecopy_context()pour hériter des contextvars dans le sous-thread.AsyncBackgroundExecutor:122— gestionnaire asynchrone, basé surasyncio.get_running_loop(), avec sémaphore optionnellemax_concurrency.gated:214—async def gated(semaphore, coro):acquire()puisrelease(), applique le plafond de concurrence autour de la coro de task.
Flux de données
_call est le cœur de l'appel de task, traduisant « appeler task_fn(arg) dans la fonction entrypoint » en « ordonnancer une tâche PUSH dans le 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_taskTrois branches décident tour à tour : (1) une task de même id est déjà en cours — réutiliser directement le future ; (2) la task a déjà tourné (writes non vide) — récupérer la valeur RETURN depuis les writes et construire un future déjà complété ; (3) task neuve — submit() remet run_with_retry au BackgroundExecutor ou AsyncBackgroundExecutor, et injecte le _call récursif comme CONFIG_KEY_CALL dans la config de la sous-task, de sorte que les appels task() internes à la sous-task empruntent la même orchestration.
Le cœur de PregelRunner.tick est une boucle FuturesDict + concurrent.futures.wait(FIRST_COMPLETED)(wait loop:281). À chaque future terminé, on dépile la tâche correspondante, on appelle commit pour écrire ses writes dans le loop, et si la task a un nœud de gestion d'erreur configuré, on soumet une task error handler. La boucle ne se termine que lorsque tous les futures en vol sont complétés (ou qu'il ne reste que le future placeholder get_waiter), rendant le contrôle à loop.after_tick pour écrire dans les canaux.
Limites et échecs
- Appel de task async en contexte synchrone échoue :
_calllèveif inspect.iscoroutinefunction(func): raise RuntimeErrorau début(sync async mismatch:716) — une entrypoint synchrone ne peut pasawait, une task asynchrone ne peut pas s'exécuter. - Exception de sous-task remonte au père :
SKIP_RERAISE_SET.add(fut)marque le future de la sous-task comme « ne pas re-raiser à la fin du tick » ; l'exception est chaînée viachainFuturevers le future du père, et le task père ne la perçoit qu'à l'await(skip reraise:781)`. Ainsi le chemin de propagation de l'erreur suit la hiérarchie d'appels, plutôt que d'être directement levé par le runner en interrompant tout le tick. __next_tick__en synchrone est différé :BackgroundExecutor.submitutilisenext_tick(ctx.run, fn, ...)(next_tick sync:66) pour placer la tâche au tick suivant ; dans l'executor asynchrone,__next_tick__est un noop, car les tasks asyncio sont naturellement ordonnancées de façon asynchrone.max_concurrencyn'affecte pas la concurrence au niveau nœud : la sémaphore deAsyncBackgroundExecutorne contraint que la coro de task (gated), pas le nombre de nœuds ordonnancés simultanément parPregelRunner.tick— la concurrence de nœuds est décidée parprepare_next_tasks, la sémaphore ne fait que throttler(semaphore:153`).- Sous-task non ré-exécutée lors d'un retry du père : si le père est retenté,
_callpasse par la branche « dedup already ran » et renvoie directement le future existant ou le résultat depuiswrites(dedup already ran:724), garantissant l'idempotence. - Héritage des contextvars :
BackgroundExecutor.submitutilisecopy_context()(copy_context:62) pour que le sous-thread hérite des contextvars du père (RunnableConfigetc.), afin que la sous-task puisse accéder à la config.
Résumé
La concurrence de @task vient de l'orchestration des futures par PregelRunner ; les abstractions centrales sont SyncAsyncFuture qui pont sync/async, _call / _acall qui convertissent l'appel de task en tâche PUSH, et BackgroundExecutor / AsyncBackgroundExecutor qui fournissent les threads / coroutines d'exécution. Pour les détails d'ordonnancement des tâches PUSH voir /subgraph/send-command, pour la sémantique de task dans entrypoint voir /func/entrypoint-task. Voir la documentation officielle : LangGraph 文档 · README