.NET · La plateforme

Parallélisme et flux

Parallel, Channel<T>, IAsyncEnumerable, annulation coopérative.

Vérifié en septembre 2026 · .NET 10 (SDK 10.0.300), C# 14 · environ 15 min

Quand un traitement a trop à faire pour un seul cœur, ou quand des éléments arrivent plus vite qu'on ne les traite, deux questions se posent : comment répartir le travail, et comment le faire circuler entre ceux qui le produisent et ceux qui le consomment. .NET y répond par quatre outils : Parallel découpe un calcul sur plusieurs threads, Channel<T> fait passer des éléments d'un producteur à un consommateur, IAsyncEnumerable<T> livre un flux élément par élément, et le jeton d'annulation arrête proprement les trois. Ce cours suppose acquis async et await (cours Programmation asynchrone), lock et Interlocked (cours Multithreading et concurrence). Toutes les sorties concurrentes citées sont des exécutions réelles, sur .NET 10 et une machine à 12 cœurs logiques ; elles varient d'une exécution à l'autre.

Parallélisme de données, concurrence de tâches

Le parallélisme de données applique la même opération à beaucoup d'éléments, en découpant la collection entre plusieurs threads qui calculent en même temps. Le gain vient des cœurs : il ne vaut que pour un travail qui occupe le processeur, et il plafonne au nombre de cœurs disponibles. C'est le domaine de Parallel.For, Parallel.ForEach et de PLINQ.

La concurrence de tâches est autre chose : plusieurs opérations indépendantes en vol en même temps, qui passent l'essentiel de leur vie à attendre un réseau ou un disque. Là, aucun cœur n'est à partager, et ajouter des threads ne sert à rien : ce sont async, Task.WhenAll et les attentes qui se recouvrent, décrits dans la section « Attendre plusieurs tâches » du cours Programmation asynchrone.

using System;
using System.Collections.Concurrent;
using System.Diagnostics;
using System.Threading;
using System.Threading.Tasks;

static class Premiers
{
    static bool EstPremier(int n)
    {
        if (n < 2) return false;

        for (var d = 2; d <= n / d; d++)
        {
            if (n % d == 0) return false;
        }

        return true;
    }

    public static void Executer()
    {
        const int Borne = 2_000_000;
        var chrono = Stopwatch.StartNew();

        var enSerie = 0;
        for (var n = 0; n < Borne; n++)
        {
            if (EstPremier(n)) enSerie++;
        }

        Console.WriteLine($"serie     {enSerie} en {chrono.ElapsedMilliseconds} ms");
        chrono.Restart();

        // Les memes donnees, decoupees en plages que plusieurs threads du pool
        // traitent en meme temps. Le corps partage deux objets : le compteur,
        // protege par Interlocked, et le dictionnaire qui releve les threads.
        // Ce dernier ne sert qu'a la demonstration, mais son TryAdd, appele
        // deux millions de fois, fait partie du temps mesure.
        var enParallele = 0;
        var appelant = Environment.CurrentManagedThreadId;
        var threads = new ConcurrentDictionary<int, bool>();

        Parallel.For(0, Borne, n =>
        {
            threads.TryAdd(Environment.CurrentManagedThreadId, true);
            if (EstPremier(n)) Interlocked.Increment(ref enParallele);
        });

        // Parallel.For est synchrone : il ne rend la main qu'une fois toutes les
        // iterations finies, et le thread appelant a travaille avec les autres.
        Console.WriteLine($"parallele {enParallele} en {chrono.ElapsedMilliseconds} ms");
        Console.WriteLine($"threads   {threads.Count}, dont l'appelant : {threads.ContainsKey(appelant)}");

        // Trois executions en Release, 12 coeurs logiques :
        // serie     148933 en 2678 ms, puis 2690 ms, puis 1690 ms
        // parallele 148933 en 779 ms, puis 476 ms, puis 421 ms
        // threads   13, dont l'appelant : True
    }
}

Sur douze cœurs logiques et trois exécutions, la version en série a pris entre 1,7 et 2,7 secondes, la version parallèle entre 0,4 et 0,8 seconde : un facteur de trois à six, pas de douze. Un cœur logique n'est pas un cœur entier : deux cœurs logiques d'un même cœur physique se partagent ses unités de calcul, beaucoup de processeurs récents mêlent cœurs rapides et cœurs d'efficacité, la fréquence baisse quand tous travaillent, et le découpage comme les objets partagés ont un coût. Parallel.For est aussi bloquant : il rend la main quand toutes les itérations sont finies, et le thread appelant y participe. Sur un serveur, ses threads viennent du même pool que ceux qui servent les requêtes ; une boucle parallèle par requête revient à faire concourir chaque requête contre toutes les autres.

Parallel.For et Parallel.ForEach

Les deux méthodes partitionnent leur source et en confient les morceaux à des tâches ; le thread appelant traite lui-même une partie des morceaux. Leur comportement se règle par ParallelOptions. La propriété qui compte est MaxDegreeOfParallelism : sa valeur par défaut, -1, ne pose aucune limite, et la boucle utilise autant de threads que l'ordonnanceur lui en fournit. La documentation signale le risque : avec des corps longs ou bloquants, le pool peut prendre la lenteur pour un manque de threads et en injecter bien plus qu'il n'en faut. Une borne explicite l'en empêche, et permet aussi de réserver une part de la machine à autre chose.

using System;
using System.Threading;
using System.Threading.Tasks;

static class Lots
{
    public static void Executer()
    {
        var mesures = new int[1_000_000];
        for (var i = 0; i < mesures.Length; i++) mesures[i] = i % 1000;

        // Quatre taches au plus, quel que soit le nombre de coeurs.
        var options = new ParallelOptions { MaxDegreeOfParallelism = 4 };
        var total = 0L;

        // Chaque tache cumule dans son propre sous-total, sans rien partager ;
        // on ne touche la variable commune qu'une fois par tache, a la fin, au
        // lieu d'une fois par element.
        Parallel.For(
            0,
            mesures.Length,
            options,
            localInit: () => 0L,
            body: (i, etat, sousTotal) => sousTotal + mesures[i],
            localFinally: sousTotal => Interlocked.Add(ref total, sousTotal));

        Console.WriteLine(total); // 499500000

        // Break arrete de distribuer les iterations au-dela de l'indice courant,
        // mais laisse finir toutes celles d'avant : l'indice le plus bas qui a
        // appele Break est donc bien le premier element qui repond au critere.
        mesures[700_000] = -1;
        mesures[300_000] = -1;

        var resultat = Parallel.For(0, mesures.Length, (i, etat) =>
        {
            if (mesures[i] < 0) etat.Break();
        });

        Console.WriteLine(resultat.IsCompleted);          // False
        Console.WriteLine(resultat.LowestBreakIteration); // 300000

        // Une exception dans le corps arrete la distribution. Celles deja levees
        // sont rassemblees : Parallel rend toujours une AggregateException,
        // meme pour une seule.
        try
        {
            Parallel.ForEach(["ventes.csv", "", "stock.csv"], fichier =>
            {
                if (fichier.Length == 0) throw new ArgumentException("nom de fichier vide");
            });
        }
        catch (AggregateException erreur)
        {
            Console.WriteLine(erreur.InnerExceptions[0].Message); // nom de fichier vide
        }
    }
}

Trois mécanismes de cet exemple reviennent souvent. La surcharge à localInit et localFinally donne à chaque tâche un état privé : le corps cumule sans synchronisation, et seule la fusion finale touche la variable commune. Le ParallelLoopState interrompt une boucle : Stop arrête au plus tôt sans rien garantir sur les indices, Break garantit que tous les indices inférieurs sont traités, ce qui rend LowestBreakIteration fiable pour une recherche du premier élément. Enfin, une exception arrête la distribution et remonte dans une AggregateException, même seule — à la différence d'un await, qui relance la première exception telle quelle.

Parallel.ForEachAsync

Parallel.ForEach est fait pour du calcul. Lui donner un travail asynchrone est une erreur classique, qui compile sans le moindre avertissement du SDK : le paramètre attendu est un Action<T>, un lambda async s'y convertit en async void, et chaque corps est compté comme terminé dès son premier await.

using System;
using System.Diagnostics;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;

static class Envois
{
    private static int enVol;
    private static int maximumEnVol;
    private static int termines;

    // L'appel au service distant : cent millisecondes d'attente pure.
    static async Task EnvoyerAsync(int commande)
    {
        var courant = Interlocked.Increment(ref enVol);
        RetenirMaximum(ref maximumEnVol, courant);

        await Task.Delay(100);

        Interlocked.Decrement(ref enVol);
        Interlocked.Increment(ref termines);
    }

    // La boucle compare-et-echange du cours Multithreading et concurrence.
    static void RetenirMaximum(ref int cible, int candidat)
    {
        var courant = Volatile.Read(ref cible);
        while (candidat > courant)
        {
            var trouve = Interlocked.CompareExchange(ref cible, candidat, courant);
            if (trouve == courant) return;
            courant = trouve;
        }
    }

    public static void Executer()
    {
        var commandes = Enumerable.Range(1, 20).ToArray();
        var options = new ParallelOptions { MaxDegreeOfParallelism = 4 };
        var chrono = Stopwatch.StartNew();

        // Parallel.ForEach attend un Action<int>. Le lambda async s'y convertit
        // sans un mot du compilateur, en async void : chaque corps rend la main
        // a son premier await, et Parallel le compte comme fini.
        Parallel.ForEach(commandes, options, async commande => await EnvoyerAsync(commande));

        Console.WriteLine($"rendu en {chrono.ElapsedMilliseconds} ms, {termines} envois termines");
        // rendu en 37 ms, 0 envois termines (97 ms a une autre execution, 0 toujours)

        Thread.Sleep(500);
        Console.WriteLine($"au plus {maximumEnVol} envois simultanes"); // au plus 20
        // La borne de 4 n'a limite que la partie synchrone de chaque corps : les
        // vingt attentes se sont superposees. Et une exception levee apres
        // l'await n'aurait aucune tache ou se poser : elle arreterait le processus.
    }
}

Trois défauts en découlent : la boucle rend la main avant qu'un seul envoi soit fini ; la borne de quatre ne limite rien, puisque les vingt attentes se superposent ; et une exception levée après l'await n'a aucune tâche où se déposer, avec les conséquences décrites pour async void dans le cours Programmation asynchrone. Depuis .NET 6, Parallel.ForEachAsync prend un corps qui rend un ValueTask et l'attend avant de passer à l'élément suivant.

using System;
using System.Diagnostics;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;

static class EnvoisBornes
{
    private static int enVol;
    private static int maximumEnVol;
    private static int termines;

    static async Task EnvoyerAsync(int commande, CancellationToken cancellationToken)
    {
        var courant = Interlocked.Increment(ref enVol);
        RetenirMaximum(ref maximumEnVol, courant);

        await Task.Delay(100, cancellationToken);

        Interlocked.Decrement(ref enVol);
        Interlocked.Increment(ref termines);
    }

    static void RetenirMaximum(ref int cible, int candidat)
    {
        var courant = Volatile.Read(ref cible);
        while (candidat > courant)
        {
            var trouve = Interlocked.CompareExchange(ref cible, candidat, courant);
            if (trouve == courant) return;
            courant = trouve;
        }
    }

    public static async Task ExecuterAsync()
    {
        var commandes = Enumerable.Range(1, 20).ToArray();
        using var source = new CancellationTokenSource(TimeSpan.FromSeconds(30));

        var options = new ParallelOptions
        {
            // Sans cette ligne, ForEachAsync prend Environment.ProcessorCount.
            MaxDegreeOfParallelism = 4,
            CancellationToken = source.Token,
        };
        var chrono = Stopwatch.StartNew();

        // Le corps rend un ValueTask : ForEachAsync l'attend avant de donner
        // l'element suivant a ce travailleur. Le jeton recu par le corps est
        // annule si l'appelant annule, ou si un autre corps echoue.
        await Parallel.ForEachAsync(commandes, options, async (commande, jeton) =>
            await EnvoyerAsync(commande, jeton));

        Console.WriteLine($"rendu en {chrono.ElapsedMilliseconds} ms, {termines} envois termines");
        // rendu en 555 ms, 20 envois termines : cinq vagues de quatre (598 ms a une
        // autre execution : Task.Delay garantit un minimum, pas une precision)
        Console.WriteLine($"au plus {maximumEnVol} envois simultanes"); // au plus 4
    }
}

MaxDegreeOfParallelism y garde la même valeur par défaut que dans les boucles synchrones, -1, mais pas le même sens : ici, -1 signifie Environment.ProcessorCount, soit douze envois simultanés au plus sur la machine de l'exemple. Pour des appels réseau, ce chiffre n'a aucun rapport avec ce que le service distant supporte, et la borne se fixe d'après lui. Quand un corps échoue, ForEachAsync cesse de distribuer et annule le jeton passé aux autres corps en cours ; l'await relance la première exception, et la propriété Exception de la tâche rendue les garde toutes, y compris les TaskCanceledException des corps annulés par ricochet : un échec et trois annulations donnent un agrégat de quatre. La méthode accepte aussi un IAsyncEnumerable<T> comme source.

Channel<T>, borné ou non

Un Channel<T> est une file sûre sous concurrence — premier entré, premier sorti, sauf le canal à priorités de CreateUnboundedPrioritized (.NET 9) —, dont les deux bouts sont séparés : un ChannelWriter<T> pour écrire, un ChannelReader<T> pour lire, et des attentes asynchrones des deux côtés. Il fait partie du framework partagé depuis .NET Core 3.0 (espace de noms System.Threading.Channels) ; le paquet NuGet du même nom le fournit aux cibles qui ne l'incluent pas, .NET Framework et .NET Standard. C'est l'équivalent asynchrone de la BlockingCollection présentée dans le cours Multithreading et concurrence.

Channel.CreateUnbounded crée un canal sans limite : chaque écriture réussit immédiatement, et si le producteur va plus vite que le consommateur, la mémoire grossit sans fin. Channel.CreateBounded fixe une capacité, et BoundedChannelFullMode dit ce qui se passe quand elle est atteinte. Une capacité nulle est admise sur .NET 10 : elle donne un canal de rendez-vous, où chaque écriture attend un lecteur.

using System;
using System.Collections.Generic;
using System.Threading.Channels;

static class ModesPlein
{
    static void Essayer(BoundedChannelFullMode mode)
    {
        var ecartes = new List<int>();

        // Trois places. La fonction de rappel recoit chaque element ecarte par
        // un mode Drop, ce qui permet au moins de le journaliser.
        var canal = Channel.CreateBounded<int>(
            new BoundedChannelOptions(3) { FullMode = mode },
            ecartes.Add);

        var refuses = new List<int>();
        for (var mesure = 1; mesure <= 5; mesure++)
        {
            // TryWrite n'attend jamais : il ecrit, ou rend false.
            if (!canal.Writer.TryWrite(mesure)) refuses.Add(mesure);
        }

        var restants = new List<int>();
        while (canal.Reader.TryRead(out var mesure)) restants.Add(mesure);

        Console.WriteLine(
            $"{mode,-10} canal [{string.Join(", ", restants)}]  " +
            $"ecartes [{string.Join(", ", ecartes)}]  refuses [{string.Join(", ", refuses)}]");
    }

    public static void Executer()
    {
        foreach (var mode in Enum.GetValues<BoundedChannelFullMode>()) Essayer(mode);
        // Wait       canal [1, 2, 3]  ecartes []  refuses [4, 5]
        // DropNewest canal [1, 2, 5]  ecartes [3, 4]  refuses []
        // DropOldest canal [3, 4, 5]  ecartes [1, 2]  refuses []
        // DropWrite  canal [1, 2, 3]  ecartes [4, 5]  refuses []

        // Un canal non borne accepte tout, tout de suite : TryWrite y rend
        // toujours true tant que le canal n'est pas clos.
        var sansBorne = Channel.CreateUnbounded<int>();
        Console.WriteLine(sansBorne.Writer.TryWrite(42)); // True
    }
}
ModeCanal pleinEmploi
Wait (défaut)WriteAsync attend une place, TryWrite rend falserien ne doit se perdre
DropOldestretire l'élément le plus ancienseules les dernières mesures comptent
DropNewestretire l'élément le plus récent du canalgarder les plus anciens en attente
DropWritejette l'élément en cours d'écriturece qui est déjà en file prime

Le piège des trois modes Drop se lit dans la sortie : TryWrite y rend toujours true, et WriteAsync s'y achève sans attendre, même quand l'élément vient d'être jeté. Le producteur n'en sait rien ; seule la fonction de rappel passée à CreateBounded voit les éléments écartés. Hors d'un cas où la perte est voulue, c'est Wait qu'il faut.

Producteur et consommateur

Avec Wait, un canal borné régule le débit de lui-même : quand il est plein, WriteAsync suspend le producteur sans bloquer de thread, jusqu'à ce qu'un consommateur libère une place. C'est la contre-pression, et c'est la raison d'être de la borne. Côté lecteur, ReadAllAsync rend un IAsyncEnumerable<T> qui se termine quand le canal est clos et vide. Le mot important est « clos » : un canal ne sait pas que son producteur a fini tant qu'on ne le lui dit pas. Le producteur ci-dessous ne le dit pas, et son consommateur attend pour toujours.

using System;
using System.Threading.Channels;
using System.Threading.Tasks;

static class PipelineOublie
{
    static async Task ProduireAsync(ChannelWriter<int> ecrivain)
    {
        for (var commande = 1; commande <= 3; commande++)
        {
            await ecrivain.WriteAsync(commande);
        }

        // Fin du travail... sans Complete. Rien dans le canal ne distingue un
        // producteur qui a fini d'un producteur qui prend son temps.
    }

    static async Task<int> ConsommerAsync(ChannelReader<int> lecteur)
    {
        var lus = 0;

        // ReadAllAsync attend l'element suivant, ou la cloture du canal.
        await foreach (var commande in lecteur.ReadAllAsync())
        {
            lus++;
        }

        return lus;
    }

    public static async Task ExecuterAsync()
    {
        var canal = Channel.CreateBounded<int>(10);

        var producteur = ProduireAsync(canal.Writer);
        var consommateur = ConsommerAsync(canal.Reader);

        await producteur;

        // Le delai n'est la que pour rendre le defaut visible : sans lui,
        // l'attente ne finirait jamais.
        try
        {
            await consommateur.WaitAsync(TimeSpan.FromSeconds(2));
        }
        catch (TimeoutException)
        {
            Console.WriteLine("le consommateur attend toujours"); // le consommateur attend toujours
        }
    }
}

Complete ferme l'écriture : les lecteurs vident ce qui reste, puis leur boucle se termine. Sa surcharge Complete(exception) fait la même chose en cas d'échec et transmet l'exception, que ReadAllAsync relance chez chaque lecteur. Clore sur tous les chemins, succès comme échec, est la règle de tout producteur.

using System;
using System.Threading;
using System.Threading.Channels;
using System.Threading.Tasks;

static class Pipeline
{
    private static int ecrits;

    static async Task ProduireAsync(ChannelWriter<int> ecrivain, CancellationToken cancellationToken)
    {
        try
        {
            for (var commande = 1; commande <= 100; commande++)
            {
                // Canal plein : WriteAsync attend qu'un consommateur libere une
                // place. Le producteur ralentit au rythme des consommateurs au
                // lieu d'empiler en memoire.
                await ecrivain.WriteAsync(commande, cancellationToken);
                Interlocked.Increment(ref ecrits);
            }

            // Plus rien ne viendra : les lecteurs finissent de vider le canal,
            // puis leur await foreach se termine.
            ecrivain.Complete();
        }
        catch (Exception erreur)
        {
            // Clore aussi sur echec : l'exception est remise aux lecteurs, que
            // ReadAllAsync relance, au lieu de les laisser attendre.
            ecrivain.Complete(erreur);
            throw;
        }
    }

    static async Task<long> ConsommerAsync(ChannelReader<int> lecteur, CancellationToken cancellationToken)
    {
        var somme = 0L;

        // Chaque element n'est lu qu'une fois : deux consommateurs se
        // partagent le flux, ils ne le recoivent pas chacun en entier.
        await foreach (var commande in lecteur.ReadAllAsync(cancellationToken))
        {
            await Task.Delay(5, cancellationToken); // le traitement, plus lent que la production
            somme += commande;
        }

        return somme;
    }

    public static async Task ExecuterAsync()
    {
        using var source = new CancellationTokenSource(TimeSpan.FromSeconds(30));

        // Dix places, et Wait, le mode par defaut, ecrit ici pour etre lu.
        var canal = Channel.CreateBounded<int>(new BoundedChannelOptions(10)
        {
            FullMode = BoundedChannelFullMode.Wait,
        });

        var producteur = ProduireAsync(canal.Writer, source.Token);
        var consommateurs = new[]
        {
            ConsommerAsync(canal.Reader, source.Token),
            ConsommerAsync(canal.Reader, source.Token),
        };

        // Juste apres le demarrage, le producteur est deja retenu : dix
        // elements dans le canal, une poignee en cours de traitement.
        await Task.Delay(20);
        Console.WriteLine($"ecrits apres 20 ms : {Volatile.Read(ref ecrits)}");
        // ecrits apres 20 ms : 16, puis 14 a une autre execution, loin des 100

        await producteur;
        var sommes = await Task.WhenAll(consommateurs);

        Console.WriteLine($"{sommes[0]} + {sommes[1]} = {sommes[0] + sommes[1]}");
        // 2523 + 2527 = 5050, puis 2579 + 2471 = 5050 : le partage varie, le total non
    }
}

Les deux consommateurs se partagent les éléments : chacun n'est lu qu'une fois, par l'un ou par l'autre, et le partage change à chaque exécution. Écrire dans un canal clos lève ChannelClosedException avec WriteAsync, rend false avec TryWrite. Les options SingleWriter et SingleReader sont des promesses de l'appelant au canal, que l'appelant doit tenir. Leur effet est plus mince qu'on ne le croit : sur .NET 10, seul un canal non borné à lecteur unique en tire une implémentation dédiée, SingleConsumerUnboundedChannel, qui perd au passage Reader.Count (CanCount y vaut false) ; un canal borné rend le même type quelles que soient les deux options.

IAsyncEnumerable et await foreach

IAsyncEnumerable<T> est l'IEnumerable<T> des flux qui attendent entre deux éléments : une page d'API, une ligne de fichier, un message. Une méthode async qui contient yield return en produit un, et await foreach le parcourt. Le compilateur y écrit GetAsyncEnumerator, une boucle de MoveNextAsync, puis DisposeAsync dans un finally. Comme un itérateur synchrone, rien ne s'exécute à l'appel : le corps démarre au premier MoveNextAsync, et chaque parcours le relance depuis le début.

L'annulation y demande un geste de plus. Le jeton peut arriver par deux portes : en argument de la méthode, ou par WithCancellation, qui le remet à GetAsyncEnumerator — la seule porte ouverte quand on reçoit un flux déjà construit. Pour que la seconde atteigne le corps de l'itérateur, son paramètre doit porter [EnumeratorCancellation]. Sans l'attribut, le jeton de WithCancellation est reçu puis ignoré.

using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;

static class FluxSourd
{
    // Le jeton est un parametre ordinaire de l'iterateur : il ne recoit que ce
    // que l'appelant passe a l'appel, ici rien, donc CancellationToken.None.
    // warning CS8425 sur ce parametre.
    static async IAsyncEnumerable<int> LireMesuresAsync(CancellationToken cancellationToken = default)
    {
        for (var page = 1; page <= 5; page++)
        {
            await Task.Delay(100, cancellationToken); // l'appel qui ramene une page
            yield return page;
        }
    }

    public static async Task ExecuterAsync()
    {
        using var source = new CancellationTokenSource(TimeSpan.FromMilliseconds(250));

        try
        {
            // WithCancellation remet le jeton a GetAsyncEnumerator. Personne
            // dans l'iterateur ne le lit : l'annulation a 250 ms passe inapercue.
            await foreach (var page in LireMesuresAsync().WithCancellation(source.Token))
            {
                Console.WriteLine($"page {page}");
            }

            Console.WriteLine("fin normale");
        }
        catch (OperationCanceledException)
        {
            Console.WriteLine("annule");
        }
        // page 1, page 2, page 3, page 4, page 5, fin normale
    }
}

Le compilateur le signale, mais par un simple avertissement, CS8425 : « L'itérateur asynchrone 'FluxSourd.LireMesuresAsync(CancellationToken)' a un ou plusieurs paramètres de type 'CancellationToken' mais aucun d'entre eux n'est décoré avec l'attribut 'EnumeratorCancellation'… » avec le SDK 10.0.300 en français. Avec l'attribut, les deux portes mènent au même paramètre ; si les deux reçoivent un jeton, l'itérateur voit un jeton lié qui s'annule dès que l'un des deux le fait.

using System;
using System.Collections.Generic;
using System.Runtime.CompilerServices;
using System.Threading;
using System.Threading.Tasks;

static class Flux
{
    // [EnumeratorCancellation] dit au compilateur : ce parametre recoit aussi
    // le jeton passe a GetAsyncEnumerator, donc celui de WithCancellation.
    // Si l'appel et WithCancellation en passent chacun un, l'iterateur recoit
    // un jeton lie aux deux.
    static async IAsyncEnumerable<int> LireMesuresAsync(
        [EnumeratorCancellation] CancellationToken cancellationToken = default)
    {
        try
        {
            for (var page = 1; page <= 5; page++)
            {
                await Task.Delay(100, cancellationToken);
                yield return page;
            }
        }
        finally
        {
            // S'execute a la fin, sur annulation, et aussi quand l'appelant
            // sort de la boucle par break : await foreach appelle DisposeAsync.
            Console.WriteLine("connexion rendue");
        }
    }

    public static async Task ExecuterAsync()
    {
        using var source = new CancellationTokenSource(TimeSpan.FromMilliseconds(250));

        try
        {
            await foreach (var page in LireMesuresAsync().WithCancellation(source.Token))
            {
                Console.WriteLine($"page {page}");
            }
        }
        catch (OperationCanceledException)
        {
            Console.WriteLine("annule");
        }
        // page 1, page 2, connexion rendue, annule

        // Rien n'est execute a l'appel : l'iterateur demarre au premier
        // MoveNextAsync, et chaque await foreach en relance un nouveau.
        await foreach (var page in LireMesuresAsync())
        {
            if (page == 2) break;
        }
        // connexion rendue
    }
}

Le finally de l'itérateur s'exécute dans tous les cas, y compris quand l'appelant sort par break : c'est l'endroit où rendre une connexion ou fermer un fichier. Dans une bibliothèque, await foreach accepte ConfigureAwait(false) sur le flux, pour la raison donnée dans le cours Programmation asynchrone. Les flux de ce type se croisent partout : ReadAllAsync d'un canal, Task.WhenEach de .NET 9, la source de Parallel.ForEachAsync.

Annulation coopérative

Le cours Programmation asynchrone pose les bases : un jeton se fait suivre de paramètre en paramètre, ThrowIfCancellationRequested vérifie entre deux opérations, une source se libère. Tout le modèle repose sur un choix : l'annulation est une demande. Le code annulé décide où s'arrêter, et il s'arrête là où son état est cohérent.

Le cas le plus courant combine deux raisons de s'arrêter : l'appelant renonce, ou l'opération dépasse son propre délai. CreateLinkedTokenSource crée une source qui s'annule dès que l'un de ses parents l'est, et CancelAfter y ajoute le délai. Au moment d'attraper l'exception, le jeton de l'appelant dit laquelle des deux a joué : l'appelant qui renonce reçoit son annulation, un délai dépassé se signale comme une panne.

using System;
using System.Threading;
using System.Threading.Tasks;

static class Delais
{
    static async Task<string> InterrogerAsync(int millisecondes, CancellationToken cancellationToken)
    {
        await Task.Delay(millisecondes, cancellationToken);
        return "tarifs";
    }

    // Le delai propre a cette operation s'ajoute a l'annulation de l'appelant
    // sans la remplacer : la source liee s'annule des que l'une des deux le fait.
    static async Task<string> LireTarifsAsync(int millisecondes, CancellationToken cancellationToken)
    {
        using var delai = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
        delai.CancelAfter(TimeSpan.FromMilliseconds(200));

        try
        {
            return await InterrogerAsync(millisecondes, delai.Token);
        }
        catch (OperationCanceledException) when (!cancellationToken.IsCancellationRequested)
        {
            // Ce n'est pas l'appelant qui a renonce : c'est notre delai. On le
            // dit avec le bon type, au lieu de faire passer une panne pour un
            // abandon volontaire.
            throw new TimeoutException("le service des tarifs n'a pas repondu en 200 ms");
        }
    }

    public static async Task ExecuterAsync()
    {
        Console.WriteLine(await LireTarifsAsync(50, CancellationToken.None)); // tarifs

        try
        {
            await LireTarifsAsync(1000, CancellationToken.None);
        }
        catch (TimeoutException erreur)
        {
            Console.WriteLine(erreur.Message); // le service des tarifs n'a pas repondu en 200 ms
        }

        // L'appelant renonce le premier (la requete HTTP est abandonnee, par
        // exemple) : l'OperationCanceledException remonte telle quelle.
        using var appelant = new CancellationTokenSource(TimeSpan.FromMilliseconds(50));
        try
        {
            await LireTarifsAsync(1000, appelant.Token);
        }
        catch (OperationCanceledException)
        {
            Console.WriteLine($"annule par l'appelant : {appelant.IsCancellationRequested}"); // True
        }
    }
}

Le filtre when fait tout le travail : il laisse passer l'annulation voulue par l'appelant et convertit l'autre. Sans lui, un service lent se présenterait à l'appelant comme un abandon qu'il n'a jamais demandé. La source liée est libérée par using, parce qu'elle détient une inscription chez son parent. Elle est préférée ici à WaitAsync, que présente le cours Programmation asynchrone : son jeton parvient jusqu'à l'opération et l'arrête, là où WaitAsync cesse d'attendre et laisse l'opération courir jusqu'au bout.

L'autre manière d'arrêter du code, l'interrompre de l'extérieur, n'existe plus. Thread.Abort injectait une exception dans un thread là où il se trouvait : au milieu d'une mise à jour, d'un constructeur statique, avant la libération d'une ressource. Depuis .NET Core, il lève PlatformNotSupportedException, et depuis .NET 5 son appel produit l'avertissement SYSLIB0006. Pour du code tiers qui ne sait pas s'annuler, la documentation renvoie à un processus séparé, que Process.Kill peut arrêter sans corrompre le processus principal.

using System;
using System.Diagnostics;
using System.Threading;
using System.Threading.Tasks;

static class Arret
{
    // Un calcul court, un peu plus d'une milliseconde : il occupe le processeur
    // sans rendre son thread.
    static void Etape() => Thread.SpinWait(20_000);

    public static void Executer()
    {
        var calcul = new Thread(() => Thread.Sleep(Timeout.Infinite)) { IsBackground = true };
        calcul.Start();

        try
        {
            // warning SYSLIB0006 : Thread.Abort est obsolete depuis .NET 5.
            calcul.Abort();
        }
        catch (PlatformNotSupportedException erreur)
        {
            Console.WriteLine(erreur.Message);
            // Thread abort is not supported on this platform.
        }

        // La version cooperative : le calcul regarde le jeton entre deux etapes,
        // et s'arrete la ou son etat est coherent.
        var chrono = Stopwatch.StartNew();
        using var source = new CancellationTokenSource(TimeSpan.FromMilliseconds(100));
        var options = new ParallelOptions
        {
            // La borne n'est pas un detail. Le minuteur de la source a lui aussi
            // besoin d'un thread du pool pour se declencher : une boucle sans
            // borne qui les occupe tous peut le retarder de plusieurs secondes.
            MaxDegreeOfParallelism = Math.Max(1, Environment.ProcessorCount / 2),
            CancellationToken = source.Token,
        };
        var lots = 0;

        try
        {
            Parallel.For(0, 1_000, options, lot =>
            {
                // Parallel ne consulte le jeton qu'entre deux iterations : un
                // corps long le regarde lui-meme.
                for (var etape = 0; etape < 50; etape++)
                {
                    source.Token.ThrowIfCancellationRequested();
                    Etape();
                }

                Interlocked.Increment(ref lots);
            });

            Console.WriteLine($"termine avant l'annulation : {lots} lots");
        }
        catch (OperationCanceledException)
        {
            // Pas d'AggregateException ici : Parallel rend l'annulation telle
            // quelle quand c'est le jeton de ses options qui l'a demandee.
            Console.WriteLine($"annule a {chrono.ElapsedMilliseconds} ms, apres {lots} lots sur 1000");
            // annule a 142 ms, apres 5 lots sur 1000
            // Seize executions : entre 113 et 201 ms, apres 0 a 7 lots.
        }
    }
}

Parallel ne regarde le jeton de ses options qu'entre deux itérations ; un corps long le consulte lui-même. Si l'exception porte ce jeton-là, la boucle la relance telle quelle, en OperationCanceledException ; celle d'un autre jeton est traitée comme une erreur ordinaire, dans une AggregateException. La borne de l'exemple est indispensable : le minuteur d'une source d'annulation a lui aussi besoin d'un thread du pool pour se déclencher. La même boucle sans borne, avec des corps qui occupent leurs threads, a vu son annulation de 100 ms arriver après 1,5 seconde, et une fois après 4,5 ; bornée à la moitié des cœurs, l'annulation est tombée entre 113 et 201 ms sur seize exécutions, après zéro à sept lots.

Ce cours vous a servi ? Offrir un café Signaler une erreur