Les flux : futures en séquence
Rappelez-vous comment, plus tôt dans ce chapitre, dans la section “passage de messages”, nous avons utilisé le récepteur pour notre canal de communication asynchrone. La méthode asynchrone recv génère une séquence d’éléments au fil du temps. Il s’agit là d’un exemple d’un modèle bien plus général connu sous le nom de flux (NDT : stream). De nombreux concepts sont naturellement représentés sous forme de flux : des éléments devenant disponibles dans une file d’attente, des morceaux de données extraits progressivement du système de fichiers lorsque l’ensemble du jeu de données est trop grange pour rentrer dans la mémoire vive de l’ordinateur, ou encore des données arrivant sur le réseau au fil du temps. Du fait que les flux soient des futures, nous pouvons les utiliser comme n’importe quelle autre sorte de future et les combiner de manière intéressante. Par exemple, nous pouvons regrouper des évènements pour éviter de déclencher trop d’appels sur le réseau, définir des délais d’expiration pour des séquences d’opérations de longue durée, ou limiter les évènements de l’interface utilisateur pour éviter d’effectuer du travail inutile.
Nous avons vu une séquence d’éléments dans le chapitre 13, lorsque nous avions regardé le trait Iterator dans la section “Le trait Iterator et la méthode next”, mais il existe deux différences entre les itérateurs et le récepteur de canal asynchrone. La première différence concerne le temps : les itérateurs sont synchrones alors que le récepteur de canal est asynchrone. La seconde différence concerne l’API. Quand on travaille directement avec Iterator, on appelle sa méthode synchrone next. Avec le flux trpl::receiver en particulier, on a appelé à la place une méthode asynchrone recv. À part ça, ces APIs se ressemblent beaucoup, et cette similarité n’est pas qu’une coïncidence ; un flux est comme une forme asynchrone d’itération. Cependant, alors que trpl::Receiver attend spécifiquement de recevoir des messages, l’API de flux polyvalente est bien plus large : elle fournit l’élément suivant de la même manière que Iterator, mais de manière asynchrone.
La similarité entre les itérateurs et les flux en Rust implique que nous pouvons créer un flux à partir de n’importe quel itérateur. Comme pour un itérateur, nous pouvons travailler avec un flux en appelant sa méthode next et ensuite attendre la sortie, comme dans le code de l’encart 17-21, qui ne se compile pas encore.
extern crate trpl; // requis par test mdbook
fn main() {
trpl::block_on(async {
let valeurs = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
let iter = valeurs.iter().map(|n| n * 2);
let mut flux = trpl::stream_from_iter(iter);
while let Some(valeur) = flux.next().await {
println!("La valeur était : {valeur}");
}
});
}
Nous commençons par un tableau de nombres, que nous convertissons en un itérateur, sur lequel nous appelons ensuite map pour doubler toutes ses valeurs. Puis nous convertissons l’itérateur en un flux en utilisant la fonction trpl::stream_from_iter. Ensuite, nous bouclons sur les éléments dans le flux au fur et à mesure de leurs arrivées avec la boucle while let.
Malheureusement, quand nous tentons d’exécuter ce code, il ne se compile pas mais, à la place, rapporte le fait qu’il n’y a pas de méthode next de disponible :
error[E0599]: no method named `next` found for struct `tokio_stream::iter::Iter` in the current scope
--> src/main.rs:10:40
|
10 | while let Some(valeur) = flux.next().await {
| ^^^^
|
= help: items from traits can only be used if the trait is in scope
help: the following traits which provide `next` are implemented but not in scope; perhaps you want to import one of them
|
1 + use crate::trpl::StreamExt;
|
1 + use futures_util::stream::stream::StreamExt;
|
1 + use std::iter::Iterator;
|
1 + use std::str::pattern::Searcher;
|
help: there is a method `try_next` with a similar name
|
10 | while let Some(valeur) = flux.try_next().await {
| ~~~~~~~~
Comme cette sortie l’explique, la raison de l’erreur de compilation est que nous avons besoin du bon trait dans la portée afin de pouvoir utiliser la méthode next. Compte tenu de la présente discussion jusqu’ici, on pourrait raisonnablement s’attendre à ce que ce trait soit Stream, mais c’est en fait StreamExt. Abréviation de extension, l’utilisation de Ext est courante dans la communauté Rust pour étendre un trait à l’aide d’un autre.
Le trait Stream concerne une interface de plus bas niveau qui combine les traits Iterator et Future. StreamExt fournit un ensemble d’APIs de plus haut niveau venant par-dessus Stream, y compris la méthode next et d’autres méthodes utilitaires similaires à celles fournies par le trait Iterator. Stream et StreamExt ne font pas encore partie de la bibliothèque standard de Rust, mais la plupart des crates de l’écosystème utilisent des définitions similaires.
La résolution de l’erreur compilation consiste à ajouter une instruction use pour trpl::StreamExt, comme dans l’encart 17-22.
extern crate trpl; // requis par test mdbook
use trpl::StreamExt;
fn main() {
trpl::block_on(async {
let valeurs = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
// -- partie masquée ici --
let iter = valeurs.iter().map(|n| n * 2);
let mut flux = trpl::stream_from_iter(iter);
while let Some(valeur) = flux.next().await {
println!("La valeur était : {valeur}");
}
});
}
Avec tous ces éléments qui s’assemblent, ce code fonctionne comme nous le souhaitons ! De plus, maintenant que nous avons StreamExt dans la portée, nous pouvons utiliser toutes ses méthodes utilitaires, exactement comme pour les itérateurs.