Comment faire l’orchestration de ressources Python ?

En Python, j’orchestre les ressources avec TaskGroup, Semaphore, timeouts, files d’attente et backpressure. Le parallélisme I/O n’est pas le vrai sujet. Le vrai sujet, c’est d’éviter de saturer vos backends, de laisser fuir des tâches, ou de casser la prod au premier timeout.

Pourquoi partir de TaskGroup ?



Je pars de TaskGroup parce que je veux une garantie de cycle de vie, pas juste une façon propre de lancer plusieurs coroutines. TaskGroup est disponible depuis Python 3.11 dans la bibliothèque standard asyncio. Il sert à faire de la concurrency structurée, c’est-à-dire que les tâches lancées ensemble vivent ensemble, et meurent ensemble si besoin.

Concrètement, toutes les tâches créées dans le bloc async with sont terminées ou annulées avant d’en sortir. Si une tâche échoue, les autres sont annulées automatiquement. C’est ça qui m’intéresse en production. Pas le côté “syntaxe élégante”.

Dans un agrégateur de dashboards, c’est typiquement le genre de problème qui finit mal. Vous interrogez en parallèle une API de prix, une base de positions, un flux d’actualité et un modèle de risque. Si le modèle de risque plante mais que le flux d’actualité continue à tourner en tâche de fond, vous venez de créer une tâche orpheline. Elle consomme peut-être du réseau, garde une connexion ouverte, loggue des erreurs trois minutes plus tard, et personne ne comprend d’où ça vient. J’ai vu ça chez un client, le bug n’était pas dans le calcul, il était dans le cycle de vie.

import asyncio
import random


async def fetch_prices(user_id):
    # Simule un appel à une API de prix
    await asyncio.sleep(0.2)
    return {"EURUSD": 1.08, "BTC": 65000}


async def fetch_positions(user_id):
    # Simule une lecture en base de positions
    await asyncio.sleep(0.15)
    return [{"asset": "BTC", "qty": 0.4}]


async def fetch_news(user_id):
    # Simule un flux d'actualité
    await asyncio.sleep(0.1)
    return ["Marché calme", "Volatilité en baisse"]


async def compute_risk(user_id):
    # Simule un modèle de risque qui peut échouer
    await asyncio.sleep(0.25)
    if random.random() < 0.1:
        raise RuntimeError(f"Risque indisponible pour {user_id}")
    return {"score": 42}


async def build_dashboard_for_user(user_id):
    # Toutes ces tâches appartiennent au même résultat métier
    async with asyncio.TaskGroup() as tg:
        prices_task = tg.create_task(fetch_prices(user_id))
        positions_task = tg.create_task(fetch_positions(user_id))
        news_task = tg.create_task(fetch_news(user_id))
        risk_task = tg.create_task(compute_risk(user_id))

    # Ici, toutes les tâches sont terminées avec succès
    return {
        "user_id": user_id,
        "prices": prices_task.result(),
        "positions": positions_task.result(),
        "news": news_task.result(),
        "risk": risk_task.result(),
    }


async def build_dashboards_for_batch(user_ids):
    tasks_by_user = {}

    # Une tâche par utilisateur, supervisée par le TaskGroup
    async with asyncio.TaskGroup() as tg:
        for user_id in user_ids:
            tasks_by_user[user_id] = tg.create_task(
                build_dashboard_for_user(user_id)
            )

    # Ici, tout le batch est fini, ou tout a été annulé proprement
    return {
        user_id: task.result()
        for user_id, task in tasks_by_user.items()
    }


async def bad_batch(user_ids):
    # À éviter : ces tâches ne sont liées à aucun cycle de vie clair
    for user_id in user_ids:
        asyncio.create_task(build_dashboard_for_user(user_id))

    # La fonction rend la main alors que le travail continue ailleurs
    return {"status": "started"}


async def main():
    dashboards = await build_dashboards_for_batch([101, 102, 103])
    print(dashboards)


if __name__ == "__main__":
    asyncio.run(main())

Je garde create_task pour les cas où je sais exactement qui supervise la tâche. Dès que plusieurs tâches contribuent au même résultat métier, je pars de TaskGroup par défaut. C’est plus simple à relire, plus sûr à annuler, et beaucoup moins piégeux quand ça casse à 2h du matin.



Comment limiter chaque backend ?



TaskGroup organise les tâches, mais il ne limite pas la capacité des ressources. C’est juste un cadre pour lancer, attendre et annuler proprement un groupe de coroutines. La limite, elle, doit être portée par des asyncio.Semaphore séparés, dimensionnés selon la vraie capacité de chaque backend.

BackendRôleLimite concurrenteRaison
PrixRécupérer les prix de marché10Backend rapide, tolère bien la charge
PositionsLire les portefeuilles clients6Base plus sensible aux pics
NewsRécupérer les actualités4API tierce avec quota serré
RisqueCalculer un score de risque3Modèle coûteux en CPU ou GPU

J’ai vu ce problème chez des clients avec des APIs tierces qui acceptent très bien 5 appels simultanés mais deviennent instables à 20. Pas besoin de dramatiser, il faut juste mettre la limite au bon endroit.

Comment faire l’orchestration de ressources Python ?
Les threads isolent surtout les I/O bloquantes ; les processus sont réservés aux calculs CPU lourds. — Source : Medium
Comment faire l’orchestration de ressources Python ?
La queue absorbe le travail en attente tandis que les workers le traitent à un rythme maîtrisé. — Source : Python documentation
Comment faire l’orchestration de ressources Python ?
Chaque sémaphore impose une pression maximale adaptée à la ressource protégée. — Source : Python documentation
import asyncio
from contextlib import asynccontextmanager

# Sémaphores déclarés au scope module.
price_sem = asyncio.Semaphore(10)
positions_sem = asyncio.Semaphore(6)
news_sem = asyncio.Semaphore(4)
risk_sem = asyncio.Semaphore(3)

risk_current = 0
risk_peak = 0
risk_lock = asyncio.Lock()


@asynccontextmanager
async def acquire_connection(name: str, semaphore: asyncio.Semaphore):
    await semaphore.acquire()

    try:
        # Ouverture simulée de connexion.
        await asyncio.sleep(0.01)
        yield
    finally:
        # Fermeture simulée de connexion.
        await asyncio.sleep(0.005)
        semaphore.release()


async def call_backend(name: str, semaphore: asyncio.Semaphore, delay: float):
    async with acquire_connection(name, semaphore):
        await asyncio.sleep(delay)


async def call_risk(dashboard_id: int):
    global risk_current, risk_peak

    async with acquire_connection("risk", risk_sem):
        async with risk_lock:
            risk_current += 1
            risk_peak = max(risk_peak, risk_current)
            assert risk_current <= 3

        await asyncio.sleep(0.08)

        async with risk_lock:
            risk_current -= 1


async def build_dashboard(dashboard_id: int):
    async with asyncio.TaskGroup() as tg:
        tg.create_task(call_backend("price", price_sem, 0.03))
        tg.create_task(call_backend("positions", positions_sem, 0.05))
        tg.create_task(call_backend("news", news_sem, 0.06))
        tg.create_task(call_risk(dashboard_id))


async def main():
    async with asyncio.TaskGroup() as tg:
        for dashboard_id in range(30):
            tg.create_task(build_dashboard(dashboard_id))

    print(f"Pic concurrent sur risk: {risk_peak}")
    assert risk_peak <= 3


asyncio.run(main())

Même avec 30 dashboards lancés en parallèle, le backend risk ne dépasse jamais 3 appels simultanés. TaskGroup donne la structure. Le sémaphore donne la pression maximale acceptable pour chaque backend.



Comment éviter les attentes infinies ?



Une orchestration fiable doit toujours borner le temps d’attente. Sinon, une ressource lente finit par bloquer toute la chaîne, et là on se retrouve avec un dashboard vide juste parce que les news financières ont mis 8 secondes à répondre. Je l’ai vu chez un client : les prix étaient prêts, les positions aussi, mais l’écran restait en chargement à cause d’un fournisseur secondaire.

Depuis Python 3.11, asyncio.timeout est la façon standard de poser un budget temps clair autour d’un appel async. L’idée est simple : chaque backend a son propre niveau de tolérance. Les news peuvent être lentes sans casser l’expérience. Le modèle de risque, lui, doit répondre vite ou être marqué indisponible.

import asyncio

BACKEND_TIMEOUTS = {
    "prices": 1.0,
    "positions": 1.5,
    "risk": 0.8,
    "news": 3.0,
}

async def call_backend(name, timeout, coro_factory, semaphore):
    client = None
    await semaphore.acquire()

    try:
        async with asyncio.timeout(timeout):
            client = await open_client(name)
            return await coro_factory(client)

    except TimeoutError:
        # Le timeout est attendu en production, pas exceptionnel.
        return None

    except asyncio.CancelledError:
        # Une annulation doit remonter, sinon l'orchestrateur croit que tout va bien.
        raise

    finally:
        try:
            if client is not None:
                await client.close()
        finally:
            # Le slot est libéré même si l'appel est annulé ou expire.
            semaphore.release()


async def build_dashboard():
    semaphore = asyncio.Semaphore(10)

    prices_task = call_backend("prices", BACKEND_TIMEOUTS["prices"], fetch_prices, semaphore)
    positions_task = call_backend("positions", BACKEND_TIMEOUTS["positions"], fetch_positions, semaphore)
    risk_task = call_backend("risk", BACKEND_TIMEOUTS["risk"], fetch_risk_model, semaphore)
    news_task = call_backend("news", BACKEND_TIMEOUTS["news"], fetch_news, semaphore)

    prices, positions, risk, news = await asyncio.gather(
        prices_task,
        positions_task,
        risk_task,
        news_task,
    )

    return {
        "prices": prices or {},
        "positions": positions or {},
        "risk": risk or {"status": "indisponible"},
        "news": news or [],
        "degraded": news is None or risk is None,
    }

Le point important, c’est que l’annulation n’est pas un cas rare. En production, un utilisateur ferme l’onglet, un job est remplacé, Kubernetes coupe un pod, une limite de temps est atteinte. Si le code avale ça proprement, tout va bien. S’il garde une connexion ouverte ou un slot de sémaphore occupé, les problèmes arrivent plus tard, et ils sont beaucoup plus durs à diagnostiquer.

Pour un dashboard, je préfère renvoyer une réponse dégradée plutôt que rien. Afficher les prix et les positions sans les news, c’est souvent acceptable. Bloquer tout l’écran pour un contenu secondaire, ça ne l’est pas.

Les erreurs classiques que j’évite systématiquement :

  • Avaler CancelledError sans le relancer.
  • Oublier le finally et laisser une connexion ou un sémaphore bloqué.
  • Mettre un timeout global trop brutal qui coupe aussi les backends rapides.
  • Utiliser le même timeout pour tous les services, alors qu’ils n’ont pas la même criticité.


Quand ajouter une file d’attente ?



La file d’attente devient utile quand le nombre de demandes entrantes dépasse ce que les backends peuvent absorber proprement. C’est le moment où je préfère ralentir proprement plutôt que lancer 10 000 tâches d’un coup et regarder l’API, la base ou le modèle de risque commencer à tousser.

Asyncio.Queue, c’est un mécanisme simple de backpressure dans la bibliothèque standard Python. Le mot veut juste dire “pression retour” : si la queue est pleine, le producteur attend. Il ne balance pas plus de travail que ce que le système peut encaisser.

La différence est très concrète. Si je fais 10 000 create_task(), je crée 10 000 coroutines planifiées tout de suite. Si j’alimente une queue consommée par 20 workers, je garde un flux contrôlé. J’ai vu ça chez un client sur un batch de dashboards lancé toutes les minutes. Le modèle de risque tenait 5 appels concurrents, pas 500. La bonne solution n’était pas “plus de CPU”, c’était une queue et un sémaphore.

import asyncio
import random

async def call_risk_model(user_id):
    await asyncio.sleep(random.uniform(0.05, 0.2))
    return {"user_id": user_id, "risk": "low"}

async def build_dashboard(user_id, risk):
    await asyncio.sleep(random.uniform(0.02, 0.1))
    print(f"Dashboard prêt pour user_id={user_id}, risk={risk['risk']}")

async def worker(name, queue, risk_sem):
    while True:
        user_id = await queue.get()

        try:
            if user_id is None:
                print(f"{name} s'arrête proprement")
                return

            async with risk_sem:
                risk = await call_risk_model(user_id)

            await build_dashboard(user_id, risk)

        finally:
            queue.task_done()

async def main():
    user_ids = range(1, 10_001)

    queue = asyncio.Queue(maxsize=500)
    risk_sem = asyncio.Semaphore(5)
    worker_count = 20

    async with asyncio.TaskGroup() as tg:
        for i in range(worker_count):
            tg.create_task(worker(f"worker-{i}", queue, risk_sem))

        for user_id in user_ids:
            await queue.put(user_id)

        for _ in range(worker_count):
            await queue.put(None)

        await queue.join()

asyncio.run(main())

Il ne faut pas confondre Queue et Semaphore. La queue limite le flux de travail en attente. Le sémaphore limite l’accès à une ressource précise, comme le modèle de risque, une API externe ou une connexion coûteuse.

OutilRôle exact
TaskGroupSupervise un groupe de tâches async et propage proprement les erreurs.
SemaphoreLimite le nombre d’accès concurrents à une ressource précise.
QueueStocke le travail en attente et impose une backpressure quand elle est pleine.

Avec ce trio, le batch peut prendre un peu plus de temps, mais il reste sain. Et dans la vraie vie, c’est souvent exactement ce qu’on veut.



Et pour le code bloquant ?



Tout ne doit pas tourner directement dans la boucle asyncio. Je vois souvent ce piège chez des équipes qui “asyncifient” un workflow, puis gardent au milieu un vieux calcul, un appel SDK bloquant, ou une lecture fichier lente. Résultat : l’event loop se fige, plus rien n’avance, même les bonnes coroutines attendent.

Quand une fonction bloque sur de l’I/O, par exemple un appel réseau synchrone ou une lecture disque, je l’isole avec asyncio.to_thread. C’est simple, lisible, parfait pour les cas ponctuels. Si je dois contrôler finement le nombre de threads, je passe sur ThreadPoolExecutor. Pour du vrai CPU-bound, donc du calcul lourd qui consomme le processeur, je regarde plutôt ProcessPoolExecutor, avec prudence, parce qu’on paie un coût de sérialisation et de transfert entre processus.

import asyncio
import time

def calcul_risque_bloquant(client_id: str) -> float:
    # Simulation d’un calcul ou appel SDK synchrone qui bloque
    time.sleep(2)
    return 0.87

async def workflow_risque(client_id: str) -> dict:
    # La fonction bloquante part dans un thread
    # L’event loop reste disponible pour le reste du workflow
    score = await asyncio.to_thread(calcul_risque_bloquant, client_id)

    return {
        "client_id": client_id,
        "score_risque": score,
        "decision": "review" if score > 0.7 else "ok",
    }

async def main():
    resultats = await asyncio.gather(
        workflow_risque("C001"),
        workflow_risque("C002"),
        workflow_risque("C003"),
    )
    print(resultats)

asyncio.run(main())

Petit détail utile côté Python 3.14 : Executor.map accepte maintenant un paramètre buffersize. Il limite le nombre de tâches soumises dont le résultat n’a pas encore été consommé. C’est très pratique pour éviter d’empiler trop de futures en mémoire sur une grosse liste. Les techniques centrales ici restent compatibles Python 3.11+, mais buffersize demande Python 3.14+.

from concurrent.futures import ThreadPoolExecutor

def enrichir_client(client_id: str) -> dict:
    # Appel I/O bloquant, par exemple un SDK legacy
    return {"client_id": client_id, "status": "enriched"}

clients = [f"C{i:05d}" for i in range(100_000)]

with ThreadPoolExecutor(max_workers=20) as executor:
    for resultat in executor.map(enrichir_client, clients, buffersize=100):
        # Seulement 100 résultats non consommés peuvent être en attente
        print(resultat)

Ma règle simple tient bien en prod :

CasChoix
Librairie async nativeCoroutine avec await
Petit appel I/O bloquantasyncio.to_thread
Beaucoup d’I/O bloquante à piloterThreadPoolExecutor
Calcul CPU lourdProcessPoolExecutor

Le bon réflexe, c’est de protéger l’event loop. Si elle respire, votre orchestration reste fluide. Si elle bloque, tout le système ralentit.



On garde quoi pour vos workers Python ?



Pour moi, une bonne orchestration Python ne consiste pas à lancer le plus de tâches possible. C’est l’inverse. Je veux savoir quelles tâches vivent ensemble, quelles ressources sont rares, combien de temps j’accepte d’attendre, et comment je ralentis quand le système est sous pression. TaskGroup donne le cadre, Semaphore protège les backends, les timeouts évitent les blocages sales, Queue ajoute de la backpressure, et les executors isolent le code bloquant. Avec ça, vous passez d’un async qui marche en démo à un async qui tient mieux en prod. Le bénéfice est simple pour vous, moins d’incidents et des traitements plus prévisibles.



FAQ



  • Pourquoi utiliser TaskGroup plutôt que create_task seul ?
    TaskGroup donne un cadre clair. Les tâches lancées dans le bloc sont supervisées ensemble. Si l’une échoue, les autres sont annulées proprement avant la sortie du bloc. Avec create_task utilisé seul, on peut vite laisser une tâche vivre en arrière-plan sans le vouloir.
  • Est-ce que TaskGroup limite automatiquement la charge sur mes APIs ?
    Non. TaskGroup structure la concurrency, mais il ne connaît pas la capacité de vos backends. Pour limiter les appels concurrents vers une API, une base ou un modèle, j’utilise un asyncio.Semaphore dédié à chaque ressource.
  • Quelle version de Python faut-il pour ces techniques ?
    Les bases présentées ici demandent Python 3.11 ou plus récent, notamment pour asyncio.TaskGroup et asyncio.timeout. La partie liée au paramètre buffersize de Executor.map demande Python 3.14 ou plus récent.
  • Queue et Semaphore servent-ils à la même chose ?
    Non. Une Queue organise le flux de travail à traiter et permet d’ajouter de la backpressure. Un Semaphore limite l’accès concurrent à une ressource précise. Dans un vrai worker, les deux se complètent très bien.
  • Comment gérer un backend lent sans casser tout le dashboard ?
    Je pose des timeouts adaptés à chaque backend et je prévois une réponse dégradée. Par exemple, je peux afficher les prix et les positions même si le flux d’actualité expire. Le point important, c’est aussi de nettoyer les connexions et de libérer les sémaphores dans un finally.

 

 

A propos de l’auteur



Je suis Franck Scandolera, expert et formateur en Tracking avancé server-side, Analytics Engineering, automatisation No/Low Code avec n8n, intégration de l’IA en entreprise et SEO/GEO. Avec mon agence webAnalyste et l’organisme Formations Analytics, j’aide des équipes à fiabiliser leurs pipelines data, leurs automatisations et leurs architectures de collecte. J’ai accompagné des clients comme Logis Hôtel, Yelloh Village, BazarChic, la Fédération Française de Football ou Texdecor. Si vous voulez industrialiser vos workflows Python, data ou IA sans bricoler dans un coin, contactez-moi.

Retour en haut