Aller au contenu principal

Concurrence et Asyncio : responsabilité des tâches, attente et annulation

Un programme concurrent peut faire avancer une opération pendant qu’une autre attend. Une requête peut attendre sa réponse tandis que le programme traite des données déjà reçues. La concurrence désigne la progression de plusieurs travaux sur des périodes qui se chevauchent, par entrelacement ou par exécution parallèle ; le parallélisme désigne leur exécution au même instant, par exemple sur des cœurs CPU différents. Faire se chevaucher les attentes peut améliorer le débit sans accélérer un calcul pris isolément.

Avant d’introduire la concurrence, déterminer qui lance chaque travail, qui reçoit son résultat et qui l’arrête en cas d’échec. Dans l’exemple ci-dessous, tous les travaux du lot ont un cycle de vie commun.

Choisir coroutines, threads ou processus​

La documentation Python sur les coroutines et les tâches distingue l’objet coroutine de la tâche ordonnancée. Appeler une fonction async def crée un objet coroutine ; un await direct attend son résultat, tandis que créer une tâche permet de la faire progresser en concurrence avec d’autres travaux. Deux attentes de coroutines consécutives peuvent donc rester séquentielles.

Les exécuteurs de threads et de processus proposent la même interface pour soumettre un travail et récupérer son résultat. Le choix dépend d’abord de ce que le travail attend et de la façon dont ses données doivent circuler :

TravailPoint de départCoût à prendre en compte
Attente réseau ou autres entrées-sorties via une API asynchrone nativeCoroutines asyncioLes appels doivent coopérer avec la boucle d’événements et éviter les blocages prolongés
Entrées-sorties bloquantes dans une bibliothèque synchrone existanteThreadPoolExecutor avec un nombre fixe de workers ; asyncio.to_thread dans un programme asynchroneLes threads partagent la mémoire : état partagé et arrêt demandent une coordination
Calcul important en Python pur nécessitant plusieurs cœursProcessPoolExecutorLe démarrage et le transfert des données ont un coût ; fonctions, arguments et résultats doivent être sérialisables par pickle, et les processus de travail doivent pouvoir importer le module principal

La condition concernant CPython dans la documentation de threading est essentielle : dans un même interpréteur dont le GIL est actif, plusieurs threads ne peuvent pas exécuter simultanément du bytecode Python. Des extensions qui libèrent le GIL peuvent calculer en parallèle. Depuis Python 3.13, les versions free-threaded permettent de désactiver le GIL, mais l’exécution simultanée de bytecode Python dans plusieurs threads d’un même interpréteur exige qu’il soit effectivement désactivé à l’exécution. Des options d’exécution ou des extensions C incompatibles peuvent le réactiver. Le choix dépend toujours de l’environnement d’exécution, des bibliothèques et de la charge de travail.

Python 3.14 propose aussi InterpreterPoolExecutor : chaque thread de travail possède son propre interpréteur et son propre GIL, ce qui permet d’utiliser plusieurs cœurs. Les interpréteurs ne peuvent pas partager directement des objets modifiables ; les fonctions soumises, leurs arguments et leurs valeurs de retour sont transmis par sérialisation avec pickle.

Appels bloquants et ordonnancement coopératif​

Une boucle d’événements exécute une seule tâche à la fois. Lorsqu’une tâche se suspend pour attendre une opération inachevée, la boucle peut faire avancer d’autres travaux. await ne garantit pas une interruption d’ordonnancement : si l’opération attendue se termine immédiatement, la tâche peut continuer sans se suspendre.

Une lecture de fichier synchrone, time.sleep ou un calcul long dans une coroutine occupe la boucle d’événements. Ajouter async def ne change pas le comportement de ces appels. asyncio.sleep suspend la tâche courante ; asyncio.to_thread déplace une fonction bloquante dans un thread. Pour une boucle de calcul intensive, envisager des processus ou des lots de calcul plus petits, avec await asyncio.sleep(0) entre les lots pour rendre la main à la boucle. Découper le calcul ou créer davantage de tâches ne suffit pas à résoudre le blocage.

Pour lancer un programme externe, les arguments, les tubes et le nettoyage des processus enfants restent soumis aux règles de gestion des sous-processus. L’ordonnancement asynchrone fait se chevaucher les attentes ; il ne gère pas à votre place tout le cycle de vie du processus enfant.

Responsabilité, durée de vie et annulation des tâches​

Un TaskGroup regroupe le cycle de vie des travaux liés dans une même portée : toutes ses tâches doivent se terminer avant la sortie du contexte. Un échec autre qu’une annulation entraîne l’annulation des autres membres ; une fois ceux-ci terminés, le groupe lève un groupe d’exceptions. Pour KeyboardInterrupt et SystemExit, c’est l’exception d’origine qui est propagée après la fin des tâches. Enregistrer les tâches enfants avec le create_task du groupe : les tâches supplémentaires créées par un appel ordinaire à asyncio.create_task ne le rejoignent pas automatiquement. TaskGroup et asyncio.timeout sont disponibles depuis Python 3.11.

Avec un appel ordinaire à create_task, l’appelant doit conserver les références, récupérer les exceptions et prévoir d’attendre les tâches lors de l’arrêt. Par défaut, gather propage la première exception sans annuler les autres travaux. Recevoir une exception ne suffit donc pas à établir que tout le groupe s’est arrêté.

Une demande d’annulation provoque CancelledError à la prochaine occasion où elle peut être traitée. Libérer les ressources dans finally ; si l’annulation est interceptée explicitement, la propager normalement après le nettoyage. L’absorber peut perturber le fonctionnement du groupe de tâches ou du délai d’attente. Après cancel(), il faut encore attendre la fin de la tâche ; l’annulation ne défait pas une écriture déjà achevée. La coordination multi-agents étend ces questions de responsabilité aux fichiers, aux messages et à plusieurs environnements d’exécution.

L’arrêt des threads obéit à une autre limite. Le Future.cancel() d’un exécuteur ne peut pas annuler un appel déjà en cours, et result(timeout=...) limite seulement l’attente de son résultat. Annuler l’attente de to_thread ne garantit donc pas l’arrêt du thread sous-jacent. Une bibliothèque bloquante doit disposer de son propre délai ou d’un mécanisme d’arrêt coopératif.

Courses sur l’état partagé​

Même avec une seule boucle d’événements, une autre tâche peut intervenir si une mise à jour se suspend entre la lecture et la réécriture de l’état. asyncio.Lock réserve un accès exclusif entre tâches d’une même boucle. Il ne synchronise pas les threads du système d’exploitation ; utiliser pour eux les outils de synchronisation de threading.

Enregistrer ce programme complet dans race.py, puis lancer python3 race.py. Il se suspend volontairement après la lecture du compteur :

import asyncio

async def count(use_lock):
value = 0
lock = asyncio.Lock()

async def update():
nonlocal value
before = value
await asyncio.sleep(0)
value = before + 1

async def worker():
if use_lock:
async with lock:
await update()
else:
await update()

async with asyncio.TaskGroup() as group:
group.create_task(worker())
group.create_task(worker())
return value

async def main():
print(f"unlocked={await count(False)}")
print(f"locked={await count(True)}")

asyncio.run(main())

Avec Python 3.13.15 et la fabrique de tâches par défaut, la sortie obtenue est :

unlocked=1
locked=2

Sans verrou, les deux workers lisent 0 puis écrivent chacun 1 : une mise à jour est perdue. Avec le verrou, le second worker lit 1 après la fin du premier et écrit 2. Le verrou doit couvrir toute l’opération lecture–modification–écriture ; verrouiller seulement l’affectation finale ne protège pas la lecture précédente.

La suspension à l’intérieur du verrou sert ici à rendre la course visible. Dans un programme ordinaire, garder les sections critiques courtes ; on peut calculer des résultats indépendants, puis confier leur regroupement à une seule tâche. Un Event signale qu’une condition s’est produite ; un Semaphore limite le nombre de tâches pouvant entrer simultanément dans une partie du travail. Créer une tâche pour chaque entrée puis les faire attendre sur un sémaphore borne le travail actif, mais pas le nombre de tâches allouées.

Files bornées, contre-pression et délais d’attente​

L’ordre FIFO et les implémentations sont présentés dans Files. Une chaîne de traitement concurrente doit aussi déterminer ce qui se passe lorsque son producteur va trop vite. Une valeur positive de maxsize dans asyncio.Queue borne les éléments en attente ; await put attend une place lorsque la file est pleine. Ce ralentissement du producteur par la capacité de l’aval s’appelle la contre-pression (backpressure). Avec maxsize=0, la longueur de la file n’est pas limitée. Cette file ne se partage pas non plus directement entre threads.

La capacité de la file, le nombre de workers et le stockage des résultats constituent des limites distinctes. Le lot ci-dessous conserve au plus 2 éléments dans la file et 2 entre les mains des workers, soit 4 travaux admis. Le producteur peut en conserver 1 de plus en attendant de l’insérer : au plus 5 éléments sont alors en cours de transmission ou de traitement. Les marqueurs d’arrêt ne sont pas des travaux. range produit les entrées progressivement, mais le dictionnaire de résultats grandit avec leur nombre. Un flux infini demande aussi une destination de résultats dont la capacité est contrôlée, et la taille des éléments doit être limitée.

asyncio.timeout annule la tâche courante lorsque son échéance est atteinte et transforme l’annulation qu’il a provoquée en TimeoutError à l’extérieur du contexte ; une annulation externe reste un CancelledError. Un délai par élément à l’intérieur du worker mesure le traitement après le retrait de la file ; un délai global autour du groupe couvre aussi les attentes d’insertion et de vidage. Un blocage de la boucle ou le nettoyage peuvent retarder le retour au-delà de l’échéance. Un délai d’attente ne force pas l’arrêt.

Un lot borné avec gestion des échecs​

Ce programme simule l’attente avec des pauses asynchrones, puis calcule les carrés des nombres de 1 à 6. Il n’utilise ni réseau ni processus externe. L’enregistrer dans bounded.py, puis lancer python3 bounded.py avec Python 3.11 ou une version ultérieure.

ModeOpération et délaisPolitique du lot
normalChaque élément attend 0,01 seconde ; délai par élément de 0,05 seconde, délai global de 2 secondesAfficher tous les résultats
item-timeoutL’élément 3 attend 0,2 seconde ; mêmes délaisNoter son dépassement et poursuivre les autres éléments
failureL’élément 1 lève ValueError après 0,01 seconde ; les autres attendent 0,2 seconde ; mêmes délaisAnnuler le groupe sans afficher de résultats partiels
batch-timeoutChaque élément attend 0,2 seconde ; délai par élément de 0,5 seconde, délai global de 0,02 secondeArrêter le lot sans afficher de résultats partiels
import asyncio

WORKERS = 2
CAPACITY = 2
STOP = object()

async def operation(number, mode):
delay = 0.01
if mode in {"failure", "batch-timeout"}:
delay = 0.2
if mode == "failure" and number == 1:
delay = 0.01
if mode == "item-timeout" and number == 3:
delay = 0.2
await asyncio.sleep(delay)
if mode == "failure" and number == 1:
raise ValueError("invalid item 1")
return number * number

async def run_batch(mode):
queue = asyncio.Queue(maxsize=CAPACITY)
results = {}
tasks = []
closed = 0
status = "ok"
item_limit = 0.5 if mode == "batch-timeout" else 0.05
batch_limit = 0.02 if mode == "batch-timeout" else 2.0

async def produce():
for number in range(1, 7):
await queue.put(number)
await queue.join()
for _ in range(WORKERS):
await queue.put(STOP)

async def consume():
nonlocal closed
try:
while True:
number = await queue.get()
try:
if number is STOP:
return
try:
async with asyncio.timeout(item_limit) as item_deadline:
value = await operation(number, mode)
except TimeoutError:
if not item_deadline.expired():
raise
results[number] = "timeout"
else:
results[number] = value
finally:
queue.task_done()
finally:
closed += 1

try:
async with asyncio.timeout(batch_limit):
try:
async with asyncio.TaskGroup() as group:
tasks.append(group.create_task(produce()))
for _ in range(WORKERS):
tasks.append(group.create_task(consume()))
except* ValueError as errors:
status = "failed"
print(f"{mode}: {sorted(str(e) for e in errors.exceptions)}")
except TimeoutError:
status = "batch-timeout"

print(f"{mode}: status={status}, workers_closed={closed}, "
f"tasks_done={all(task.done() for task in tasks)}")
if status == "ok":
ordered = sorted(results.items())
total = sum(value for value in results.values()
if isinstance(value, int))
print(f"results={ordered}, total={total}")

async def main():
for mode in ("normal", "item-timeout", "failure", "batch-timeout"):
await run_batch(mode)

asyncio.run(main())

L’exécution avec Python 3.13.15 a produit cette sortie, les résultats étant triés par numéro d’entrée :

normal: status=ok, workers_closed=2, tasks_done=True
results=[(1, 1), (2, 4), (3, 9), (4, 16), (5, 25), (6, 36)], total=91
item-timeout: status=ok, workers_closed=2, tasks_done=True
results=[(1, 1), (2, 4), (3, 'timeout'), (4, 16), (5, 25), (6, 36)], total=82
failure: ['invalid item 1']
failure: status=failed, workers_closed=2, tasks_done=True
batch-timeout: status=batch-timeout, workers_closed=2, tasks_done=True

La somme normale vaut 1 + 4 + 9 + 16 + 25 + 36 = 91 ; après le dépassement de délai de l’élément 3, elle vaut 91 - 9 = 82. Ici, status=ok signifie que le lot a été traité selon sa politique. Le dépassement reste explicite dans les résultats ; ce statut ne signifie pas que toutes les opérations ont réussi.

run_batch affiche le statut et les résultats, puis renvoie None. Après avoir intercepté TimeoutError, le worker consulte la méthode expired() du contexte de délai pour vérifier que sa propre échéance a été atteinte. Si l’opération lève elle-même cette exception avant l’échéance, le worker la relance ; le groupe annule les autres travaux et propage l’échec à l’appelant.

Le producteur attend queue.join() avant d’envoyer un marqueur d’arrêt par worker. Chaque get réussi correspond à un seul task_done, y compris pour les marqueurs d’arrêt, les dépassements de délai et les exceptions. Un get qui n’a pas encore réussi ne décrémente pas le compteur. Ce décompte coordonne le vidage ; il ne prouve pas la réussite du travail métier.

Si un worker échoue, le groupe annule le producteur qui attend dans join ou put, ainsi que les autres workers. Ce chemin d’échec n’attend plus le vidage : les éléments non retirés sont abandonnés avec cette file locale. except* ValueError traite ValueError et ses sous-classes comme des échecs de lot dans cet exemple ; les exceptions qui ne correspondent pas remontent à l’appelant dans un groupe d’exceptions, et les annulations externes se propagent aussi. Lorsque le délai global expire, le même groupe se termine avant l’affichage du statut de dépassement.

Chaque sortie workers_closed=2 et tasks_done=True montre que les deux workers sont passés par leur nettoyage de sortie et que les trois tâches enfants sont terminées. Le regroupement intervient après la fin du groupe. Les workers écrivent sous des numéros distincts, sans await entre ces opérations : aucun verrou supplémentaire n’est nécessaire ici. Un compteur partagé avec une attente entre sa lecture et son écriture nécessiterait la synchronisation montrée plus haut.

Pour utiliser un service réel, remplacer operation par un appel asynchrone annulable tout en conservant la responsabilité des tâches et les deux périmètres de délai. Cet exemple ne réessaie pas les opérations. Avant d’ajouter des tentatives, fixer le délai total, leur nombre maximal et l’idempotence des effets de bord : l’annulation du groupe ne défait pas une opération distante déjà achevée.

Explorer les liensOuvrir le réseau