Загружаем научный разбор
Подготавливаем текст, источники и редакционные примечания без изменения разметки страницы.
Подготавливаем текст, источники и редакционные примечания без изменения разметки страницы.
Автор: Казачкин Даниил Михайлович · Обновлено
Почему async/await не ограничивает память автоматически: модель очереди, Web Streams и собственный пример согласования производителя с потребителем.
Приложение получает большой файл и отправляет его порциями в медленное хранилище. Чтение занимает миллисекунды, запись каждой порции — значительно дольше. Если производитель продолжает добавлять работу без ограничения, данные копятся в памяти. Отсутствие блокировки главного потока ещё не означает, что программа контролирует потребление ресурсов.
Backpressure — обратный сигнал, сообщающий производителю, что потребитель не готов принимать данные с прежней скоростью. Это часть протокола взаимодействия, а не просто задержка между запросами. Замедление должно следовать из состояния очереди и завершения работы, иначе выбранная пауза окажется случайной настройкой для одной машины.
В собственном численном сценарии производитель создаёт 200 блоков в секунду, а потребитель успевает обработать 80. При устойчивой разнице очередь растёт примерно на 120 блоков каждую секунду. Если блок занимает 64 КиБ, только полезные данные очереди добавляют около 7,5 МиБ в секунду. Здесь не учтены объекты, копии и сетевые буферы, поэтому это модель причины роста, а не замер процесса.
Конечная память заставляет выбрать политику: замедлить источник, остановить приём, отклонить часть работы либо перенести буфер в ограниченное внешнее хранилище. Бесконечная очередь просто откладывает момент отказа. Для задачи, где нельзя терять данные, производитель должен уметь ждать или корректно завершать операцию при невозможности продолжения.
WHATWG Streams задаёт очереди, стратегии их размера и распространение backpressure через потоковые конвейеры. High water mark участвует в вычислении желаемого размера очереди, но не является универсальным запретом выделить больше памяти: источник может проигнорировать сигнал, а отдельный блок оказаться огромным. Документация Node.js описывает реализацию Web Streams API для серверного JavaScript.
Важно различать очередь блоков и очередь байтов. Лимит в четыре блока может означать четыре маленькие строки или четыре многомегабайтных массива. Если продукт ограничивает память, размер блока и функция оценки очереди должны соответствовать этому требованию. Кроме того, учитывайте буферы до источника и после приёмника: одна аккуратная очередь не описывает всю систему.
Код предназначен для Node.js с Web Streams API. Источник выдаёт шесть маленьких чисел по запросу pull. Приёмник искусственно задерживает завершение каждой записи. Мы проверяем порядок и число одновременно выполняющихся записей; максимальный размер всего процесса этот опыт не измеряет.
const { ReadableStream, WritableStream } = require('node:stream/web');
const assert = require('node:assert/strict');
async function run() {
let produced = 0;
let active = 0;
let maxActive = 0;
const received = [];
const input = new ReadableStream({
pull(controller) {
if (produced === 6) controller.close();
else controller.enqueue(produced++);
},
}, { highWaterMark: 1 });
const output = new WritableStream({
async write(value) {
active++;
maxActive = Math.max(maxActive, active);
await new Promise(resolve => setTimeout(resolve, 2));
received.push(value);
active--;
},
}, { highWaterMark: 1 });
await input.pipeTo(output);
assert.deepEqual(received, [0, 1, 2, 3, 4, 5]);
assert.equal(maxActive, 1);
console.log('completed:', received.length);
}
run().catch(error => { console.error(error); process.exitCode = 1; });Обещание, возвращаемое write, завершается после обработки блока. Если вместо await запустить запись и немедленно вернуть управление, поток сочтёт блок обработанным раньше времени. Тогда собственная скрытая очередь асинхронных операций может расти, хотя внешняя потоковая очередь выглядит свободной. Поэтому место завершения Promise является частью контракта ресурса.
Аналогичная ошибка возникает с массивом: Promise.all(items.map(save)) запускает все save сразу. Конструкция ожидает завершения, но не ограничивает число начатых операций. Для независимых задач может понадобиться отдельный ограничитель параллелизма; для упорядоченного потока — согласованный потоковый конвейер. Эти требования следует различать до выбора удобного синтаксиса.
Кроме нормальной передачи нужно проверить отказ потребителя, отмену пользователя и внезапное окончание источника. Определите, кто закрывает файл, освобождает соединение и прекращает получение новых блоков. При ошибке нельзя бесконечно продолжать чтение данных, которым уже некуда отправляться. Если имеются внешние ресурсы, их очистка должна быть привязана к реальному жизненному циклу операции.
Для диагностики записывайте глубину очереди, средний размер блока, скорость приёма и обработки, время ожидания и число активных операций. Снижение CPU при растущей очереди не обязательно означает запас производительности: узкое место может находиться в диске, сети или внешнем сервисе. Измерение каждого этапа помогает локализовать причину.
Учитывайте момент выделения блока. Если источник сначала прочитал весь файл в большой массив и лишь затем выдаёт из него небольшие срезы, аккуратная передача не отменит уже занятую память. Для ограниченного потребления исходное чтение тоже должно быть порционным. Аналогично потребитель, сохраняющий каждый полученный блок в массив ради будущего объединения, переносит накопление за пределы потоковой очереди. Проверьте обе границы отдельно и опишите, какие данные обязаны оставаться живыми до окончания операции. Это помогает отличать необходимость задачи от случайного удержания объектов.
Стандартизованный поток не гарантирует низкую память при любых пользовательских callbacks. Наш пример подтверждает порядок и последовательную обработку маленьких значений, но не заменяет нагрузочный тест экспорта гигабайтного файла. Backpressure полезен в загрузках, архивировании, обработке логов и выгрузках базы. Его практическая ценность — заранее определённое поведение при медленном потребителе, которое можно проверить без ожидания аварийного заполнения памяти.
Самостоятельный русскоязычный разбор ЯдроКода. Описания первоисточников отделены от авторских учебных примеров и инженерных выводов. Материал не является переводом или перепечаткой.
Учебные данные, расчёты, таблицы и программные примеры созданы для этой публикации. Иллюстрации и программный код из первоисточников не воспроизводятся.