Сравните семантику Future и Stream: почему одно ожидание Future получает итоговое значение, тогда как Stream нужно опрашивать многократно?
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.
Здесь каждый вызов next().await получает максимум один элемент. Сам stream завершает последовательность только после того, как выдаст все значения и вернёт None.
Главный практический компромисс — управление давлением. Последовательное чтение и обработка элемента позволяют не загружать память всеми событиями сразу, но медленный потребитель может задержать дальнейшее чтение. Буферизация или параллельная обработка повышают пропускную способность, однако требуют ограничения размера очереди и контроля порядка результатов.
Сервис получает события из сетевого источника и должен обрабатывать их по одному, не загружая в память весь журнал. Возможны три варианта: собрать всё в коллекцию, создавать отдельную задачу на каждый элемент или представить источник как stream.
Сбор всех данных упрощает код, но нарушает потоковую обработку и может привести к исчерпанию памяти. Отдельная задача на событие даёт параллелизм, но без ограничения числа задач создаёт неконтролируемую нагрузку и усложняет завершение.
Выбран Stream с ограниченной буферизацией и последовательным next().await. Это сохраняет потоковую обработку и естественное обратное давление; параллелизм можно добавить отдельным ограниченным этапом, если обработка элементов независима. В результате память остаётся предсказуемой, а завершение определяется единым сигналом None.
Можно ли вернуть Option<T> из Future и считать его Stream?
Нет. Future<Output = Option<T>> возвращает одно значение типа Option<T> и завершается после этого. Если результатом является Some(value), это всё равно один результат; None не означает возможность получить следующий элемент.
Stream же может многократно возвращать Some(value) в отдельных циклах опроса. Его None означает завершение последовательности, а не просто особое значение единственного результата.
Почему Stream не обязан выдавать элементы при каждом опросе?
Как и future, stream может вернуть Pending. Например, сетевой источник ещё не получил пакет, поэтому stream сохраняет состояние операции и предоставляет исполнителю Waker.
После уведомления исполнитель снова опросит stream. Это позволяет не блокировать поток runtime и не тратить процессор на активное ожидание.
Что произойдёт, если потребитель stream будет медленнее производителя?
Всё зависит от реализации источника. Stream может читать данные только по запросу потребителя, использовать ограниченный буфер и тем самым создавать обратное давление. Либо он может иметь внутреннюю очередь, которая будет расти и приведёт к увеличению задержек и расхода памяти.
Поэтому при проектировании нужно отдельно определить размер буфера, политику переполнения и возможность отмены. Сам факт, что источник представлен как Stream, не гарантирует ограниченное потребление памяти или автоматическую параллельную обработку.