Программирование RustКонкурентность и asyncРазработчик Rust, создающий асинхронные сервисы

Сравните семантику Future и Stream: почему одно ожидание Future получает итоговое значение, тогда как Strea...

Сравните семантику Future и Stream: почему одно ожидание Future получает итоговое значение, тогда как Stream нужно опрашивать многократно?

Проходите собеседования с ИИ помощником Hintsage

Краткий ответ

Future описывает получение одного результата: после Poll::Ready операция завершена. Stream описывает последовательность результатов: каждый успешный опрос выдаёт очередной элемент, а завершение обозначается None, поэтому для чтения всей последовательности нужны многократные ожидания.

Исторический контекст

Асинхронные операции бывают одноразовыми и многократными. Разделение этих моделей позволяет исполнителю одинаково эффективно работать как с единичным ответом сервиса, так и с потоком событий, не создавая отдельную задачу для каждого элемента.

В Rust базовая модель Future предназначена для одного результата. Абстракция Stream обычно предоставляется async-экосистемой, например crate futures, и добавляет к идее future возможность выдавать значения порциями.

Постановка проблемы

Если представить поток событий как один Future, придётся самостоятельно кодировать состояние между событиями, признак завершения и повторную регистрацию пробуждения. Ошибка в таком коде может привести к пропущенному событию, преждевременному завершению или бесконечному ожиданию.

Обратная ошибка также опасна: запуск отдельной async-задачи на каждый элемент усложняет ограничение числа одновременно обрабатываемых событий и может создать неограниченный рост очереди задач.

Подробное решение

При опросе Future возможны два основных результата: Pending, если результат ещё не готов, и Ready(value), когда единственный результат получен. После Ready future считается завершённым, поэтому корректный исполнитель не должен продолжать обычный жизненный цикл её опроса.

У Stream аналогичная модель применяется к каждому элементу. Опрос может вернуть Pending, Ready(Some(item)) для очередного значения или Ready(None) для конца последовательности. Между выдачами элементов stream сохраняет своё состояние и при необходимости регистрирует Waker для следующего события.

Метод next() обычно является адаптером: он возвращает отдельный future, который ожидает один следующий элемент stream. Поэтому цикл с повторным next().await фактически последовательно опрашивает stream, пока не получит None.

use futures::{stream, StreamExt}; #[tokio::main] async fn main() { let mut events = stream::iter([10, 20, 30]); while let Some(event) = events.next().await { println!("{event}"); } }

Здесь каждый вызов next().await получает максимум один элемент. Сам stream завершает последовательность только после того, как выдаст все значения и вернёт None.

Главный практический компромисс — управление давлением. Последовательное чтение и обработка элемента позволяют не загружать память всеми событиями сразу, но медленный потребитель может задержать дальнейшее чтение. Буферизация или параллельная обработка повышают пропускную способность, однако требуют ограничения размера очереди и контроля порядка результатов.

Ситуация из практики

Сервис получает события из сетевого источника и должен обрабатывать их по одному, не загружая в память весь журнал. Возможны три варианта: собрать всё в коллекцию, создавать отдельную задачу на каждый элемент или представить источник как stream.

Сбор всех данных упрощает код, но нарушает потоковую обработку и может привести к исчерпанию памяти. Отдельная задача на событие даёт параллелизм, но без ограничения числа задач создаёт неконтролируемую нагрузку и усложняет завершение.

Выбран Stream с ограниченной буферизацией и последовательным next().await. Это сохраняет потоковую обработку и естественное обратное давление; параллелизм можно добавить отдельным ограниченным этапом, если обработка элементов независима. В результате память остаётся предсказуемой, а завершение определяется единым сигналом None.

Что кандидаты часто упускают

  1. Можно ли вернуть Option<T> из Future и считать его Stream?

    Нет. Future<Output = Option<T>> возвращает одно значение типа Option<T> и завершается после этого. Если результатом является Some(value), это всё равно один результат; None не означает возможность получить следующий элемент.

    Stream же может многократно возвращать Some(value) в отдельных циклах опроса. Его None означает завершение последовательности, а не просто особое значение единственного результата.

  2. Почему Stream не обязан выдавать элементы при каждом опросе?

    Как и future, stream может вернуть Pending. Например, сетевой источник ещё не получил пакет, поэтому stream сохраняет состояние операции и предоставляет исполнителю Waker.

    После уведомления исполнитель снова опросит stream. Это позволяет не блокировать поток runtime и не тратить процессор на активное ожидание.

  3. Что произойдёт, если потребитель stream будет медленнее производителя?

    Всё зависит от реализации источника. Stream может читать данные только по запросу потребителя, использовать ограниченный буфер и тем самым создавать обратное давление. Либо он может иметь внутреннюю очередь, которая будет расти и приведёт к увеличению задержек и расхода памяти.

    Поэтому при проектировании нужно отдельно определить размер буфера, политику переполнения и возможность отмены. Сам факт, что источник представлен как Stream, не гарантирует ограниченное потребление памяти или автоматическую параллельную обработку.