Если производитель AsyncStream стабильно быстрее потребителя при неограниченной буферизации, к чему это приведёт?
При неограниченной буферизации производитель не будет ждать потребителя: элементы продолжат накапливаться в памяти. Если разрыв скоростей сохраняется, очередь может расти до исчерпания памяти; такая схема не создаёт обратного давления.
AsyncStream появился как структурированный способ представить callback-источники, события и делегаты в форме AsyncSequence. Исходная проблема заключалась в необходимости безопасно передавать поток значений из событийной API в асинхронный потребитель, не управляя вручную состоянием ожидания и завершения.
Буферизация нужна потому, что производитель и потребитель не обязаны работать одновременно. Однако выбор неограниченного буфера переносит ответственность за разницу скоростей на потребление памяти.
Представим телеметрию или события датчика: производитель публикует тысячи значений в секунду, а потребитель обрабатывает их медленнее. При неограниченной буферизации ни одно значение не отбрасывается автоматически, поэтому отставание превращается в растущую очередь.
Последствия — повышенное потребление памяти, задержка обработки устаревших данных и в крайнем случае аварийное завершение из-за нехватки памяти. При этом сам факт успешной постановки элемента в буфер не означает, что потребитель уже его обработал.
Политика bufferingPolicy определяет поведение, когда потребитель не успевает за производителем. Политика unbounded сохраняет все элементы и не ограничивает размер очереди, поэтому она приемлема только при обоснованном ограничении объёма или гарантированно сопоставимых скоростях.
Ограниченные политики меняют семантику данных. bufferingNewest сохраняет самые свежие элементы и удаляет более старые при переполнении; это подходит для состояния, где важнее актуальное значение. bufferingOldest сохраняет уже находящиеся в очереди элементы и отбрасывает новые, что полезно, когда порядок и ранее полученные события важнее свежести.
Ограниченная буферизация не является полноценным обратным давлением: производитель обычно не блокируется, а часть значений теряется. Если нельзя терять элементы, нужно отдельно проектировать механизм регулирования скорости — например, подтверждения, ограниченную очередь с ожиданием или другой протокол между производителем и потребителем.
Здесь в буфере хранится максимум одно наиболее новое значение. Быстрый производитель не ждёт медленного потребителя, но при переполнении старые значения могут быть отброшены; результат yield позволяет обнаружить факт отбрасывания.
Сервис мониторинга получает частые обновления загрузки процессора, но экран обновляется только несколько раз в секунду. Использование unbounded сохраняло бы все промежуточные значения, увеличивало задержку и расход памяти, хотя пользователю нужны только свежие данные.
Вариант с bufferingOldest сохранял бы устаревшие значения и делал экран ещё менее актуальным. Вариант с неограниченной очередью не терял бы данные, но требовал бы отдельной гарантии, что потребитель догонит производителя.
Выбором становится bufferingNewest(1): старые промежуточные обновления можно отбросить, а потребитель получает наиболее свежее состояние. Для финансовых операций или событий аудита такой выбор был бы неверен — там нужна доставка каждого события и явное управление скоростью, а не потеря данных.
Нет, сама по себе она обычно не обеспечивает блокирующее обратное давление. Производитель продолжает вызывать yield, а политика решает, какой элемент удалить или отклонить. Если ожидание потребителя обязательно, его нужно реализовать отдельным протоколом.
Нет. Успешный результат означает, что значение принято stream-континуейшном — например, помещено в буфер или передано ожидающему потребителю. Обработка телом цикла for await может произойти позже, поэтому для подтверждения обработки нужен отдельный механизм квитирования.
BufferingNewest подходит для быстро меняющегося состояния: координат, прогресса или показателей, где устаревшие значения бесполезны. BufferingOldest подходит, когда важнее сохранить уже поступившие элементы, но допустимо потерять новые при перегрузке. Если недопустима потеря любого элемента, обе политики не подходят без дополнительного контроля скорости и ёмкости очереди.