Skip to content

Exécution concurrente : orchestration de futures par @task

源码版本1.2.9

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_options récupère la callback CONFIG_KEY_CALL depuis la config (c'est le _call / _acall injecté par le runner), emballe la task en une structure Call, et l'insère dans le canal TASKS via schedule_task(schedule PUSH task:713) ; prepare_next_tasks au tour suivant la trouvera naturellement.
  • Abstraction unifiée sync / async : SyncAsyncFuture(SyncAsyncFuture:253) est à la fois un concurrent.futures.Future et un objet avec __await__ — son __await__ fait un seul yield, 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, _call cherche d'abord dans futures un future de même id, et le réutilise s'il existe ; si la task a déjà fini (next_task.writes non vide), il récupère directement la valeur RETURN depuis 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, AsyncBackgroundExecutor lit config["max_concurrency"] et instancie un asyncio.Semaphore ; toutes les coros de task sont emballées dans gated(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étient submit / put_writes / node_error_handler_map / schedule_error_handler.
  • PregelRunner.tick:176 — exécute un lot de tâches en synchrone ; boucle concurrent.futures.wait(FIRST_COMPLETED) pour récupérer les futures terminés et commit les writes.
  • PregelRunner.atick:360 — version asynchrone, asyncio.wait remplace le wait synchrone, sémantique d'itérateur identique.
  • _call:700 — callback CONFIG_KEY_CALL synchrone, 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 — via schedule_task (loop.accept_push / aaccept_push) insère la structure Call dans le canal TASKS.
  • SyncAsyncFuture:253 — objet pont à la fois future synchrone et awaitable asynchrone ; __await__ fait un seul yield pour rendre le contrôle.
  • _call_with_options:276 — implémentation de _TaskFunction.__call__, récupère CONFIG_KEY_CALL depuis la config et l'invoque.
  • BackgroundExecutor:40 — gestionnaire de pool de threads synchrones, submit utilise copy_context() pour hériter des contextvars dans le sous-thread.
  • AsyncBackgroundExecutor:122 — gestionnaire asynchrone, basé sur asyncio.get_running_loop(), avec sémaphore optionnelle max_concurrency.
  • gated:214async def gated(semaphore, coro) : acquire() puis release(), 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 » :

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

Trois 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 : _call lève if inspect.iscoroutinefunction(func): raise RuntimeError au début(sync async mismatch:716) — une entrypoint synchrone ne peut pas await, 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 via chainFuture vers 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.submit utilise next_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_concurrency n'affecte pas la concurrence au niveau nœud : la sémaphore de AsyncBackgroundExecutor ne contraint que la coro de task (gated), pas le nombre de nœuds ordonnancés simultanément par PregelRunner.tick — la concurrence de nœuds est décidée par prepare_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é, _call passe par la branche « dedup already ran » et renvoie directement le future existant ou le résultat depuis writes(dedup already ran:724), garantissant l'idempotence.
  • Héritage des contextvars : BackgroundExecutor.submit utilise copy_context()(copy_context:62) pour que le sous-thread hérite des contextvars du père (RunnableConfig etc.), 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