Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Mise en œuvre de la concurrence avec Async

Dans cette partie, nous allons appliquer la programmation asynchrone à certains des mêmes défis liés à la concurrence que nous avons abordés avec les tâches au chapitre 16. Comme nous avons déjà abordé bon nombre des concepts clés à ce sujet, nous nous concentrerons ici sur les différences entre les tâches et les futures.

Dans bien des cas, les APIs permettant de gérer laconcurrence en utilisant la programmation asynchrone sont très similaires à celles pour l’utilisation des tâches. Dans d’autres cas, elles s’avèrent assez différentes. Même quand les APIs semblent similaires entre tâches et l’asynchronisme, elles ont souvent des comportements differents — et leurs caractéristiques de performances sont presque toujours différentes.

Créer une nouvelle tâche avec spawn

La première opération que nous avons abordée dans la section “Créer une nouvelle tâche avec spawn du chapitre 16 consistait à faire un comptage sur deux tâches distinctes. Faisons de même en utilisant la programmation asynchrone. La crate trpl fournit une fonction spawn_task qui semble très similaire à l’API thread::spawn API, ainsi qu’une fonction sleep qui est une version asynchrone de l’API thread::sleep. Nous pouvons les utiliser ensemble pour implémenter l’exemple de comptage, comme le montre l’encart 17-6.

Filename: src/main.rs
extern crate trpl; // requis par test mdbook

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        trpl::spawn_task(async {
            for i in 1..10 {
                println!("coucou numéro {i} depuis la première tâche !");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        });

        for i in 1..5 {
            println!("coucou numéro {i} depuis la seconde tâche !");
            trpl::sleep(Duration::from_millis(500)).await;
        }
    });
}
Listing 17-6: Creating a new task to print one thing while the main task prints something else

Comme point de départ, nous configurons notre fonction main avec trpl::block_on de manière à ce que notre fonction de plus haut niveau puisse être asynchrone.

Note : dans ce chapitre, à partir d’ici, chaque exemple inclura exactement le même code d’encapsulation avec trpl::block_on dans main ; nous passerons donc souvent tout cela sous silence, tout comme nous le faisons avec main. N’oubliez pas d’inclure cela dans votre code !

Nous écrivons ensuite deux boucle dans chaque bloc, chacune contenant un appel à trpl::sleep, laquelle attend une demi-seconde (500 millisecondes) avant d’envoyer le prochain message. Nous mettons une boucle dans le corps d’un trpl::spawn_task et l’autre dans une boucle for de plus haut niveau. Nous ajoutons aussi un await après les appels à sleep.

Ce code se comporte de manière similaire à l’implémentation basée sur les tâches — y compris le fait que vous pourriez voir les messages apparaître dans un ordre différent dans votre propre terminal quand vous l’exécutez :

coucou numéro 1 depuis la seconde tâche !
coucou numéro 1 depuis la première tâche !
coucou numéro 2 depuis la première tâche !
coucou numéro 2 depuis la seconde tâche !
coucou numéro 3 depuis la première tâche !
coucou numéro 3 depuis la seconde tâche !
coucou numéro 4 depuis la première tâche !
coucou numéro 4 depuis la seconde tâche !
coucou numéro 5 depuis la première tâche !

Cette version s’arrête dès que la boucle for dans le corps du bloc asynchrone principal se termine, car la tâche lancée par spawn_task est arrêtée lorsque la fonction main prend fin. Si vous voulez qu’elle s’exécuter jusqu’au bout de la tâche, vous devrez utiliser un sélecteur de jointure (NDT : join handle) pour attendre l’achèvement de la première tâche. Avec les tâches, nous avons utilisé la méthode join pour “bloquer” jusqu’à ce que la tâche ait fini son exécution. Dans l’encart 17-7, nous pouvons utiliser await pour faire de même, car le sélecteur de jointure est lui-même une future. Son type Output est un Result, donc nous le déballons également après l’avoir attendu.

Filename: src/main.rs
extern crate trpl; // requis par test mdbook

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let selecteur = trpl::spawn_task(async {
            for i in 1..10 {
                println!("coucou numéro {i} depuis la première tâche !");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        });

        for i in 1..5 {
            println!("coucou numéro {i} depuis la seconde tâche !");
            trpl::sleep(Duration::from_millis(500)).await;
        }

        selecteur.await.unwrap();
    });
}
Listing 17-7: Using await with a join handle to run a task to completion

Cette version mise à jour s’exécute jusqu’à ce que les deux boucles se terminent :

coucou numéro 1 depuis la seconde tâche !
coucou numéro 1 depuis la première tâche !
coucou numéro 2 depuis la première tâche !
coucou numéro 2 depuis la seconde tâche !
coucou numéro 3 depuis la première tâche !
coucou numéro 3 depuis la seconde tâche !
coucou numéro 4 depuis la première tâche !
coucou numéro 4 depuis la seconde tâche !
coucou numéro 5 depuis la première tâche !
coucou numéro 6 depuis la première tâche !
coucou numéro 7 depuis la première tâche !
coucou numéro 8 depuis la première tâche !
coucou numéro 9 depuis la première tâche !

Jusqu’ici, on dirait que la programmation asynchrone et la programmation multitâches donnent les mêmes résultats, avec juste des syntaxes différentes : utiliser await au lieu d’appeler join sur le sélecteur de jointure, et attendre les appels sleep.

La plus grande différence est que nous n’avons pas eu besoin de lancer une autre tâche du système d’exploitation pour faire ceci. En fait, nous n’avons même pas besoin de lancer une tâche ici. Du fait que les blocs asynchrones se compilent en des futures anonymes, nous pouvons mettre chaque boucle dans un bloc asynchrone et avoir le moteur d’exécution qui les fait tourner toutes les deux jusqu’à leur achèvement, en utilisant la fonction trpl::join.

Dans la section “Attendre que toutes les tâches aient fini” du chapitre 16, nous vous avons montré comment utiliser la méthode join sur le type JoinHandle renvoyé quand vous appelez std::thread::spawn. La fonction trpl::join est similaire, mais pour les futures. Quand vous lui donnez deux futures, elle produit une seule nouvelle future dont la sortie est un tuple contenant la sortie de chaque future passée en argument, une fois qu’elles se sont achevées _toutes les deux. Ainsi, dans l’encart 17-8, nous utilisons trpl::join pour attendre fut1 ainsi and que fut2 s’achèvent. Nous _n’_attendons pas fut1 et fut2, mais à la place la nouvelle future produite par trpl::join. Nous ignorons la sortie, car c’est simplement un tuple contenant deux valeurs unitaires.

Filename: src/main.rs
extern crate trpl; // requis par test mdbook

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let fut1 = async {
            for i in 1..10 {
                println!("coucou numéro {i} depuis la première tâche !");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let fut2 = async {
            for i in 1..5 {
                println!("coucou numéro {i} depuis la seconde tâche !");
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        trpl::join(fut1, fut2).await;
    });
}
Listing 17-8: Using trpl::join to await two anonymous futures

Quand nous exécutons ceci, nous voyons que les deux futures sont exécutées jusqu’à leur terme :

coucou numéro 1 depuis la première tâche !
coucou numéro 1 depuis la seconde tâche !
coucou numéro 2 depuis la première tâche !
coucou numéro 2 depuis la seconde tâche !
coucou numéro 3 depuis la première tâche !
coucou numéro 3 depuis la seconde tâche !
coucou numéro 4 depuis la première tâche !
coucou numéro 4 depuis la seconde tâche !
coucou numéro 5 depuis la première tâche !
coucou numéro 6 depuis la première tâche !
coucou numéro 7 depuis la première tâche !
coucou numéro 8 depuis la première tâche !
coucou numéro 9 depuis la première tâche !

Maintenant, vous allez voir exactement le même ordre à chaque fois, ce qui diffère beaucoup de ce qu’on a vu avec les tâches et avec trpl::spawn_task dans l’encart in 17-7. C’est parce que la fonction trpl::join est équitable (NDT : fair), dans le sens où elle vérifie chaque future à intervalles réguliers, en alternant entre elles, et ne laisse jamais l’une prendre de l’avance si l’autre est prête. En programmation multitâches, le système d’exploitation décide quelle tâche choisir et combien de temps on la laisse s’exécuter. En Rust asynchrone, c’est le moteur d’exécution qui décide quelle tâche choisir (en pratique, les détails deviennent plus compliqués car un moteur d’exécution asynchrone peut utiliser le multitâches du système d’exploitation sous le capot dans la manière dont il gère la concurrence, ce qui fait que garantir l’équité peut nécessiter plus de travail au moteur d’exécution — mais cela reste possible !). Les moteurs d’exécution ne sont pas tenus de garantir l’équité pour une opération donnée, ils proposent souvent différentes APIs pour vous laisser décider si vous souhaitez ou non l’équité.

Essayez quelques de ces variantes d’attente des futures et voyez ce qu’elles font :

  • enlevez le bloc asynchrone de l’entourage d’une des boucles ou des deux ;
  • attendez chaque bloc asynchrone immédiatement après l’avoir défini ;
  • entourez la première boucle uniquement dans un bloc asynchrone, puis attendez la future résultante après le corps de la seconde boucle.

À titre de défi supplémentaire, voyez si vous pouvez déterminer, dans chaque cas, quelle sera la sortie avant d’exécuter le code !

Échange de données entre deux tâches en utilisant le passage de messages

Échanger des données entre des futures vous semblera également familier : nous utiliserons derechef le passage de messages, mais cette fois-ci avec les versions asynchrones des types et des fonctions. Nous emprunterons un cheminement légèrement différent de celui suivi dans la section “Transfert de données entre fils d’exécution avec l’envoi de messages” du chapitre 16 afin d’illustrer certaines des différences clés entre la concurrence basée sur les tâches et celle basée sur les futures. Dans l’encart 17-9, nous allons commencer par un seul bloc asynchrone — sans créer un travail distinct, comme nous avions créé un d’un fil d’exécution (tâche) séparé.

Filename: src/main.rs
extern crate trpl; // requis par test mdbook

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let valeur = String::from("salut");
        tx.send(val).unwrap();

        let reception = rx.recv().await.unwrap();
        println!("reçu '{reception}'");
    });
}
Listing 17-9: Creating an async channel and assigning the two halves to tx and rx

Nous utilisons ici trpl::channel, une version asynchrone de l’API du canal de communication avec plusieurs producteurs et un seul consommateur (NDT :“multiple producer, single consumer) que nous avions utilisée avec le multitâches dans le chapitre 16. La version asynchrone de l’API est juste un tout petit peu différente de la version multitâches : elle utilise un récepteur mutable au lieu d’un récepteur rx immuable, et sa méthode recv génère une future que nous devons attendre, au lieu de générer directement la valeur. Nous pouvons maintenant envoyer des messages vers le récepteur. Vous remarquerez que nous n’avons pas besoin de lancer un fil d’exécution séparé ni même une tâche ; nous avons tout juste à attendre l’appel rx.recv.

La méthode synchrone Receiver::recv de std::mpsc::channel bloque jusqu’à ce qu’elle reçoive un message. La méthode trpl::Receiver::recv ne bloque pas, car elle est asynchrone. Au lieu de bloquer, elle rend le contrôle au moteur d’exécution jusqu’à ce qu’un message soit reçu ou que la partie émettrice du canal se ferme. En revanche, nous n’attendons pas l’appel send, car il ne bloque pas. Il n’en a pas besoin, car le canal vers lequel nous l’envoyons est illimité.

Note : comme tout ce code asynchrone s’exécute dans un bloc asynchrone dans un appel à trpl::block_on, tout ce qui se trouve à l’intérieur peut éviter le blocage. Cependant, le code situé en-dehors bloquera jusqu’à ce que la fonction block_on renvoie une valeur. C’est là tout l’intérêt de la fonction trpl::block_on : elle vous permet de choisir où bloquer sur un ensemble de code asynchrone, et donc où faire la transition entre le code synchrone et le code asynchrone.

Retenez deux choses à propos de cet exemple. D’abord, le message arrivera tout de suite. Ensuite, bien que nous utilisions ici une future, il n’y a pas encore de concurrence. Tout dans le programme se déroule de manière séquentielle, exactement comme s’il n’y avait pas de futures impliquées.

Occupons-nous de la première partie en envoyant une série de messages et en faisant une pause entre eux, comme montré dans l’encart 17-10.

Filename: src/main.rs
extern crate trpl; // requis par test mdbook

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let valeurs = vec![
            String::from("salut"),
            String::from("à partir"),
            String::from("de la"),
            String::from("future"),
        ];

        for valeur in valeurs {
            tx.send(valeur).unwrap();
            trpl::sleep(Duration::from_millis(500)).await;
        }

        while let Some(valeur) = rx.recv().await {
            println!("reçu '{valeur}'");
        }
    });
}
Listing 17-10: Sending and receiving multiple messages over the async channel and sleeping with an await between each message

En plus d’envoyer des messages, nous devons pouvoir les réceptionner. Dans ce cas, comme nous savons combien de messages sont en train d’arriver, nous pourrions opérer manuellement en appelant rx.recv().await quatre fois. Dans la pratique, toutefois, nous serons généralement amenés à attendre un nombre indéterminé de messages ; nous devons donc continuer à attendre jusqu’à déterminer qu’il ne subsiste plus de messages.

Dans l’encart 16-10, nous avons eu recours à une boucle for pour traiter tous les éléments reçus à partir d’un canal de communication synchrone. Rust n’a toutefois pas encore la possibilité d’utiliser une boucle for avec une série d’éléments produits de manière asynchrone, nous devons donc utiliser une boucle que nous n’avons pas encore vue : la boucle conditionnelle while let. Il s’agit de la version en boucle de la construction if let vue dans la section “Gestion concise du flux d’éxécution avec if let et let...else du chapitre 6. La boucle poursuivra son exécution tant que le motif qu’elle comporte continue à correspondre avec la valeur.

L’appel à rx.recv produit une future, que nous attendons. Le moteur d’exécution va mettre la future en attente jusqu’à ce qu’elle soit prête. Dès qu’un message arrive, la future va se résoudre en Some(message) autant de fois qu’il y a d’arrivées de message. Lorsque le canal se referme, peu importe qu’il y ait eu ou non des messages, la future va sera alors résolue en None afin d’indiquer qu’il n’y a plus de valeurs et que nous devons donc cesser l’interrogation — c’est-à-dire arrêter d’attendre.

La boucle while let rassemble tous ces éléments. Si le résultat de l’appel à rx.recv().await est Some(message), nous avons accès au message et nous pouvons l’utiliser dans le corps de la boucle, exactement comme avec if let. Si le résultat est None, la boucle se termine. À chaque fois que la boucle s’achève, elle atteint de nouveau le point d’attente, donc le moteur d’exécution fait de nouveau une pause jusqu’à l’arrivée d’un autre message.

Le code parvient désormais à envoyer et recevoir tous les messages. Malheureusement, il reste encore quelques problèmes. D’une part, les messages n’arrivent pas avec des intervalles d’une demi-seconde : ils arrivent tous en même temps, deux secondes (2000 millisecondes) après le lancement du programme. D’autre part, ce programme ne se termine jamais ! Au lieu de cela, il attend pour l’éternité de nouveaux messages. Vous devez l’arrêter à l’aide de ctrl-C.

Le code dans un bloc asynchrone s’exécute de manière linéaire

Commençons par regarder pourquoi les messages arrivent tous en un seul coup après le délai complet, au lieu d’arriver les uns après les autres avec des intervalles entre chacun. À l’intérieur d’un bloc asynchrone donné, l’ordre dans lequel le mot-clé await apparaît est également l’ordre dans lequel ils sont exécutés lors de l’exécution du programme.

Il n’y a qu’un seul bloc asynchrone dans l’encart 17-10, donc tout son contenu s’exécute de manière linéaire ; il n’y a toujours pas de concurrence. Tous les appels tx.send ont lieu, entrecoupés par tous les appels trpl::sleep et de leurs points d’attente associés. Ça n’est qu’après cela que la boucle while let peut passer par un des points d’attente await sur les appels recv.

Afin d’avoir le comportement souhaité, à savoir quand le délai d’attente a lieu entre chaque message, nous devons placer les opérations tx et rx dans leurs propres blocs asynchrones, comme montré dans l’encart 17-11. Alors, le moteur d’exécution peut exécuter chacun d’eux séparément en utilisant trpl::join, comme dans l’encart 17-8. Une fois encore, nous attendons le résultat de l’appel à trpl::join et non pas les futures individuelles. Si nous attendions les futures individuelles les unes après les autres, cela reviendrait simplement à un flux séquentiel, soit exactement ce que nous essayons de ne pas faire.

Filename: src/main.rs
extern crate trpl; // requis par test mdbook

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx_fut = async {
            let vals = vec![
                String::from("salut"),
                String::from("à partir"),
                String::from("de la"),
                String::from("future"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let rx_fut = async {
            while let Some(valeur) = rx.recv().await {
                println!("reçu '{valeur}'");
            }
        };

        trpl::join(tx_fut, rx_fut).await;
    });
}
Listing 17-11: Separating send and recv into their own async blocks and awaiting the futures for those blocks

Avec le code mis à jour de l’encart 17-11, les messages sont affichés à un intervalle de 500 millisecondes, au lieu de tous en même temps après 2 secondes.

Déplacement de la possession dans un bloc asynchrone

Le programme ne se termine toujours pas, du fait de la manière dont la boucle while let interagit avec trpl::join :

  • La future renvoyée depuis trpl::join ne s’achève qu’une fois que toutes les deux futures qui lui ont été transmises se sont terminées.
  • La future tx_fut se termine une fois qu’elle a fini sa période d’attente après l’envoi du dernier message de vals.
  • La future rx_fut ne se terminera pas avant que la boucle while let ne se finisse.
  • La boucle while let ne se terminera pas tant que l’attente de rx.recv ne renverra pas None.
  • L’attente de rx.recv renverra None uniquement une fois que l’autre bout du canal de communication est refermé.
  • Le canal de communication ne se fermera que si nous appelons rx.close ou si le côté émetteur, tx, est interrompu.
  • Nous n’appelons rx.close nulle part, et tx ne sera pas interrompu jusqu’à ce que le bloc asynchrone le plus exterme transmis à trpl::block_on se termine.
  • Le bloc ne peut pas se terminer car il reste bloqué en attendant que trpl::join soit terminé, ce qui nous ramène au début de cette liste.

Pour le moment, le bloc asynchrone à partir duquel nous envoyons les messages ne fait qu’emprunter tx parce qu’envoyer un message n’a pas besoin d’avoir la propriété, mais si nous pouvions déplacer tx à l’intérieur de ce bloc asynchrone, il serait libéré une fois ce bloc terminé. Dans la section “Capture des références ou transfert de propriété” du chapitre 16, nous avons souvent besoin de déplacer des données dans des fermetures lorsque nous travaillons avec des tâches. La même dynamique de base s’applique aux blocs asynchrones, de sorte que le mot-clé move fonctionne avec les blocs asynchrones de la même manière qu’avec les fermetures.

Dans l’encart 17-12, nous modifions le bloc utilisé pour envoyer les messages de async en async move.

Filename: src/main.rs
extern crate trpl; // requis par test mdbook

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx_fut = async move {
            // -- partie masquée ici --
            let vals = vec![
                String::from("salut"),
                String::from("à partir"),
                String::from("de la"),
                String::from("future"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let rx_fut = async {
            while let Some(valeur) = rx.recv().await {
                println!("reçu '{valeur}'");
            }
        };

        trpl::join(tx_fut, rx_fut).await;
    });
}
Listing 17-12: A revision of the code from Listing 17-11 that correctly shuts down when complete

Quand nous exécutons cette version du code, il se termine correctement après que le dernier message a été émis et reçu. Maintenant, voyons ce que nous devrions changer pour émettre des données à partir de plus d’une future.

Combiner plusieurs futures avec la macro join!

Ce canal de communication asynchrone est également un canal à producteurs multiples, de sorte que nous pouvons appeler clone sur tx si nous voulons envoyer des messages à partir de plusieurs futures, comme il est montré dans l’encart 17-13.

Filename: src/main.rs
extern crate trpl; // requis par test mdbook

use std::time::Duration;

fn main() {
    trpl::block_on(async {
        let (tx, mut rx) = trpl::channel();

        let tx1 = tx.clone();
        let tx1_fut = async move {
            let vals = vec![
                String::from("salut"),
                String::from("à partir"),
                String::from("de la"),
                String::from("future"),
            ];

            for val in vals {
                tx1.send(val).unwrap();
                trpl::sleep(Duration::from_millis(500)).await;
            }
        };

        let rx_fut = async {
            while let Some(valeur) = rx.recv().await {
                println!("reçu '{valeur}'");
            }
        };

        let tx_fut = async move {
            let vals = vec![
                String::from("plus de"),
                String::from("messages"),
                String::from("pour"),
                String::from("vous"),
            ];

            for val in vals {
                tx.send(val).unwrap();
                trpl::sleep(Duration::from_millis(1500)).await;
            }
        };

        trpl::join!(tx1_fut, tx_fut, rx_fut);
    });
}
Listing 17-13: Using multiple producers with async blocks

Tout d’abord, nous clonons tx, créant tx1 en-dehors du premier bloc asynchrone. Nous déplaçons tx1 dans ce bloc, exactement comme nous l’avions fait précédemment avec tx. Puis, plus tard, nous déplaçons le tx original dans un nouveau bloc asynchrone, où nous envoyons plus de messages avec un délai un petit peu plus lent. Il se trouve que nous mettons ce nouveau bloc asynchrone après le bloc asynchrone qui reçoit les messages, mais il aurait tout aussi bien pu se trouver avant. L’essentiel réside dans l’ordre dans lequel les futures sont attendues, et non pas celui dans lequel elles sont créées.

Ces deux blocs asynchrones pour émettre les messages doivent être des blocs async move, de telle sorte que tx et tx1 soient tous deux supprimés une fois que ces blocs se terminent. Si ça n’était pas le cas, nous nous retrouverions dans la même boucle infinie que celle dans laquelle nous nous sommes retrouvés au départ.

Finalement, nous passons de trpl::join à trpl::join! afin de pouvoir gérer la future supplémentaire : la macro join! attend un nombre arbitraire de futures alors que nous connaissons le nombre de futures au moment de la compilation. Nous aborderons plus loin dans ce chapitre l’attente d’un ensemble de futures dont le nombre est inconnu.

Nous voyons maintenant tous les messages provenant des futures émettrices, et du fait que ces futures d’émission utilisent des délais légèrement différents après l’émission, les messages sont aussi reçus à ces différents intervalles :

reçu 'salut'
reçu 'plus de'
reçu 'à partir'
reçu 'de la'
reçu 'messages'
reçu 'future'
reçu 'pour'
reçu 'vous'

Nous avons vu comment utiliser le passage de messages pour envoyer des données entre des futures, comment le code à l’intérieur d’un bloc asynchrone s’exécute de manière séquentielle, comment transférer la possession vers un bloc asynchrone, et comment combiner plusieurs futures. Voyons ensuite comment et pourquoi indiquer au moteur d’exécution qu’il peut passer à une autre tâche.