Streams: futures в последовательности
Вспомните, как ранее в этой главе, в разделе «Передача
сообщений», мы использовали приемник нашего
async-канала. Async-метод recv со временем производит последовательность
элементов. Это пример гораздо более общего шаблона, известного как stream.
Многие концепции естественно представляются в виде streams: элементы,
становящиеся доступными в очереди; фрагменты данных, постепенно извлекаемые из
файловой системы, когда полный набор данных слишком велик для памяти
компьютера; или данные, поступающие по сети с течением времени. Поскольку
streams являются futures, мы можем использовать их с любыми другими видами
futures и комбинировать интересными способами. Например, можно объединять
события в пакеты, чтобы не запускать слишком много сетевых вызовов, задавать
тайм-ауты для последовательностей долгих операций или ограничивать частоту
событий пользовательского интерфейса, чтобы не выполнять лишнюю работу.
Мы уже видели последовательность элементов в главе 13, когда рассматривали
трейт Iterator в разделе «Трейт Iterator и метод next», но между итераторами и приемником async-канала есть два различия.
Первое различие — время: итераторы синхронны, а приемник канала асинхронен.
Второе различие — API. Работая напрямую с Iterator, мы вызываем его
синхронный метод next. В частности, со stream trpl::Receiver вместо этого
мы вызывали асинхронный метод recv. В остальном эти API ощущаются очень
похожими, и это сходство не случайно. Stream похож на асинхронную форму
итерации. Однако если trpl::Receiver конкретно ожидает получения сообщений,
то stream API общего назначения гораздо шире: он предоставляет следующий
элемент так же, как это делает Iterator, но асинхронно.
Сходство между итераторами и streams в Rust означает, что мы на самом деле
можем создать stream из любого итератора. Как и с итератором, со stream можно
работать, вызывая его метод next, а затем ожидая результат, как в листинге
17-21, который пока не скомпилируется.
extern crate trpl; // required for mdbook test
fn main() {
trpl::block_on(async {
let values = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
let iter = values.iter().map(|n| n * 2);
let mut stream = trpl::stream_from_iter(iter);
while let Some(value) = stream.next().await {
println!("The value was: {value}");
}
});
}
Мы начинаем с массива чисел, который преобразуем в итератор, а затем вызываем
у него map, чтобы удвоить все значения. Затем мы преобразуем итератор в
stream с помощью функции trpl::stream_from_iter. Далее мы проходим по
элементам stream по мере их поступления с помощью цикла while let.
К сожалению, когда мы пытаемся запустить код, он не компилируется, а вместо
этого сообщает, что доступного метода next нет:
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(value) = stream.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(value) = stream.try_next().await {
| ~~~~~~~~
Как объясняет этот вывод, причина ошибки компилятора в том, что для
использования метода next нужный трейт должен находиться в области видимости.
С учетом нашего обсуждения до этого момента можно было бы разумно ожидать, что
этим трейтом будет Stream, но на самом деле это StreamExt. Сокращение
Ext от extension — распространенный в сообществе Rust шаблон для
расширения одного трейта другим.
Трейт Stream определяет низкоуровневый интерфейс, который фактически
объединяет трейты Iterator и Future. StreamExt предоставляет более
высокоуровневый набор API поверх Stream, включая метод next, а также
другие вспомогательные методы, похожие на те, которые предоставляет трейт
Iterator. Stream и StreamExt пока не являются частью стандартной
библиотеки Rust, но большинство крейтов экосистемы используют похожие
определения.
Исправление ошибки компилятора состоит в том, чтобы добавить объявление use
для trpl::StreamExt, как в листинге 17-22.
extern crate trpl; // required for mdbook test
use trpl::StreamExt;
fn main() {
trpl::block_on(async {
let values = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
// --snip--
let iter = values.iter().map(|n| n * 2);
let mut stream = trpl::stream_from_iter(iter);
while let Some(value) = stream.next().await {
println!("The value was: {value}");
}
});
}
Когда все эти части собраны вместе, этот код работает так, как нам нужно!
Более того, теперь, когда StreamExt находится в области видимости, мы можем
использовать все его вспомогательные методы, как и с итераторами.