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

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, который пока не скомпилируется.

Имя файла: src/main.rs
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}");
        }
    });
}
Listing 17-21: Создание stream из итератора и печать его значений

Мы начинаем с массива чисел, который преобразуем в итератор, а затем вызываем у него 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.

Имя файла: src/main.rs
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}");
        }
    });
}
Listing 17-22: Успешное использование итератора как основы для stream

Когда все эти части собраны вместе, этот код работает так, как нам нужно! Более того, теперь, когда StreamExt находится в области видимости, мы можем использовать все его вспомогательные методы, как и с итераторами.