Асинхронные итераторы и генераторы: потоки по запросу без буфера
Асинхронные итераторы работают по запросу: потребитель тянет каждое значение, поэтому продюсер приостановлен на yield, а backpressure структурный. break/throw запускает finally генератора, и курсоры закрываются. Цепочка генераторов — ленивый конвейер с памятью O(1).
Твой ночной импортёр работал нормально месяцами, потом клиент с выгрузкой в 9 миллионов строк запустил его, и под умер на 1.5 ГБ с JavaScript heap out of memory. Ты свёл это к одной невинной строке: const rows = await db.query('SELECT * FROM events'); rows.map(transform). Драйвер забуферил весь результирующий набор в JS-массив до того, как твой код увидел хоть одну строку, поэтому пиковая память масштабировалась с размером выгрузки, а не с размером твоей пачки. С трансформом всё было в порядке; баг был в том, что ты затянул всё в память разом вместо строки за раз. Фикс — не под побольше, а смена того, кто задаёт темп.
Протокол асинхронной итерации: next() возвращает промис
Зачем это знать инженеру? Именно протокол даёт for await гарантию ограниченной памяти — понимание его позволяет строить дата-конвейеры, которые никогда не упадут в OOM независимо от размера результирующего набора, и объяснить ревьюеру, почему твой код безопасен там, где наивный await db.query(...) — нет.
Обычный итерируемый объект предоставляет Symbol.iterator; асинхронный предоставляет Symbol.asyncIterator — метод, который возвращает итератор, чей next() возвращает Promise<{ value, done }> вместо простого { value, done }. Именно это единственное отличие — результат является промисом — позволяет каждому шагу что-то await-ить (сетевой round-trip, чтение с диска) перед тем, как выдать следующее значение. for await (const x of it) — чистый сахар: он один раз вызывает it[Symbol.asyncIterator](), затем крутит цикл const { value, done } = await iter.next() пока done не станет true, выполняя тело для каждого value. Он последователен по построению — цикл полностью await-ит один next() перед тем, как просить следующий.
Протокол почти никогда не пишут руками. Асинхронный генератор — async function* с yield — это эргономичный инструмент автора: он строит соответствующий протоколу асинхронный итератор за тебя, и каждый yield приостанавливает функцию до того, как потребитель снова вызовет next(). Функция буквально останавливается на yield и не движется дальше, пока её не потянут.
// Пагинируем cursor-based API, выдавая по одной записи за раз.
async function* fetchAllUsers(client) {
let cursor = null;
do {
const page = await client.get("/users", { cursor, limit: 100 });
for (const user of page.items) {
yield user; // suspends here until the consumer pulls the next value
}
cursor = page.nextCursor;
} while (cursor); // only fetches the next page when the buffer is drained
}
// Потребляем: по одному пользователю, не более одной страницы в памяти.
for await (const user of fetchAllUsers(client)) {
await index(user); // each await fully completes before the next user is pulled
}Цикл по странице не может убежать вперёд потребителя: выдав последнего пользователя страницы N, генератор заморожен посреди for; он возобновляется — и только тогда отправляет запрос за страницей N+1 — когда index(user) зарезолвится и for await снова вызовет next(). Память ограничена одной страницей плюс одной записью в полёте, есть у API 100 пользователей или 100 миллионов.
Pull побеждает push: backpressure бесплатно
Это и есть senior-инсайт. Колбэки, EventEmitter-ы и readable.on('data', …) — это push: продюсер решает, когда данные приходят, и швыряет их в тебя. Если ты не успеваешь — медленная запись в БД за быстрым чтением из сокета — события всё равно идут, и ты обязан реализовать backpressure руками: приостановить источник, буферизовать, отслеживать размер буфера, возобновить при опустошении. Забудь что-нибудь из этого — и неограниченный буфер превращается в тот самый OOM из вступления, просто переехавший.
Асинхронные итераторы — это pull: ничего не производится, пока потребитель не вызовет next(). Поэтому backpressure структурный, а не приделанная функция — продюсер приостановлен на yield и физически неспособен убежать вперёд, потому что единственное, что его возобновляет, — это твой следующий запрос. Медленный потребитель просто тянет медленнее; продюсер ждёт. А поскольку Node Readable реализует Symbol.asyncIterator, ты получаешь это бесплатно поверх реального I/O:
import { createReadStream } from "node:fs";
import { createInterface } from "node:readline";
// readline over a Readable is itself async-iterable.
const rl = createInterface({
input: createReadStream("huge.log"), // a 40 GB file is fine
crlfDelay: Infinity,
});
for await (const line of rl) {
await ship(line); // if ship() is slow, the file read naturally throttles
}Пока ship(line) в ожидании, for await не вызывает next(), поэтому readline не запрашивает больше, поэтому нижележащий файловый поток не читается дальше своего highWaterMark. Чтение ОС приостанавливается. Ты не написал ни .pause(), ни .resume(), ни учёта буфера — pull-модель это обеспечила.
▸Почему это работает
Почему push — модель сложнее, если обе двигают одни и те же байты? В push полномочие на тайминг живёт у продюсера, а ограничение по ёмкости — у потребителя; они расщеплены через границу, поэтому корректность зависит от канала обратной связи (pause/resume), который ты обязан подключить и никогда не уронить. В pull и полномочие, и ограничение живут у потребителя: одна сторона решает и когда, и стоит ли продвигаться, поэтому нет канала, который можно было бы неверно настроить. Backpressure перестаёт быть протоколом и становится инвариантом самого потока управления.
Очистка при раннем выходе: курсоры закрываются в finally
Pull-продюсер, удерживающий ресурс — курсор БД, открытый файловый дескриптор, соединение из пула, — несёт реальный риск: потребитель может сделать break, return или throw из цикла for await до того, как источник исчерпан. Рантайм это обрабатывает. Когда ты покидаешь for await рано, он вызывает метод return() итератора. Внутри асинхронного генератора return() заставляет приостановленный yield повести себя так, будто там выполнился return: код после текущего yield не выполняется, но любой охватывающий блок finally выполняется. Этот finally — единственное надёжное место, чтобы освободить ресурс.
async function* streamRows(pool, sql) {
const conn = await pool.connect();
const cursor = conn.cursor(sql);
try {
for (let batch; (batch = await cursor.read(500)).length; ) {
for (const row of batch) yield row; // consumer may break out mid-stream
}
} finally {
// Выполняется при нормальном исчерпании И при break/return/throw через return() генератора.
await cursor.close();
conn.release(); // the leak you'd get without finally: a held connection forever
}
}
for await (const row of streamRows(pool, "SELECT * FROM events")) {
if (await isDuplicate(row)) break; // triggers return() → the finally above runs
await upsert(row);
}Режим отказа конкретен: помести cursor.close() и conn.release() после цикла вместо finally — и первый же break их пропустит. Соединение никогда не вернётся в пул; сделай так несколько тысяч раз — пул исчерпан, и каждый следующий запрос виснет. Одна тонкость, которую стоит держать на радаре: return() срабатывает, только если итерация реально началась и затем остановилась через break/return/throw или доработала до конца. Если ты захватишь итератор вручную (const it = gen[Symbol.asyncIterator]()), вытянешь пару значений, а затем просто сбросишь ссылку, не вызвав it.return() и не исчерпав его, генератор останется приостановленным на своём последнем yield навсегда, и finally никогда не выполнится — ресурс утекает. for await всегда делает правильно; ручная итерация делает очистку твоей задачей.
Композиция: ленивые конвейеры при постоянной памяти
Поскольку генератор и является асинхронно-итерируемым, и может потреблять такой, ты можешь сцепить их в конвейер, где каждая стадия трансформирует по одному элементу за раз. Сравни жадный стиль с массивами — (await getAll()).filter(...).map(...), который материализует каждый промежуточный массив, с ленивым конвейером из генераторов, который держит по одному элементу на стадию:
async function* mapGen(src, fn) { for await (const x of src) yield fn(x); }
async function* filterGen(src, ok) { for await (const x of src) if (ok(x)) yield x; }
// source → filter → map → consumer, one row flowing through at a time
const pipeline =
mapGen(
filterGen(streamRows(pool, "SELECT * FROM events"), (r) => r.amount > 0),
(r) => ({ id: r.id, cents: Math.round(r.amount * 100) }),
);
for await (const record of pipeline) {
await sink.write(record); // backpressure propagates all the way to the cursor
}Арифметика памяти — в этом весь смысл. Материализация 10М строк в среднем по ~150 байт — это ~1.5 ГБ живого массива — ровно тот OOM из вступления, и становится хуже с каждым промежуточным .map/.filter, который выделяет ещё один полный массив. Конвейер из генераторов — O(1) по числу строк: в любой момент есть одна строка в streamRows, одна в filterGen, одна в mapGen, одна записывается — четыре объекта, а не десять миллионов. Ты можешь стримить результирующий набор больше, чем RAM. А поскольку pull-backpressure прошивает каждую стадию, медленный sink.write придушит курсор БД на дальнем конце без единой лишней строки кода.
Нужно трансформировать и записать 10М строк БД в воркере, с ограниченной пиковой памятью, и хочется, чтобы медленный приёмник придушивал чтение. Какой подход подходит?
Потребитель делает `for await (const row of streamRows(...)) { ...; break; }` после нескольких строк. У генератора есть try/finally, закрывающий курсор. Что выполнится?
- 01Почему backpressure асинхронного итератора «структурный» и бесплатный, а push-based EventEmitter заставляет реализовывать его руками?
- 02Где закрывать курсор БД в асинхронном генераторе и что именно утекает, если этого не делать?
Асинхронно-итерируемый объект предоставляет Symbol.asyncIterator, чей next() возвращает Promise<{ value, done }>; for await — сахар для одного его вызова и цикла await-нутых next() пока не done, последовательно. Асинхронные генераторы (async function* + yield) пишут их эргономично: каждый yield приостанавливает функцию, пока потребитель снова не потянет, поэтому продюсер не может убежать вперёд. Это и есть ключевое отличие от push (колбэки, EventEmitter, .on('data')), где продюсер задаёт темп, а ты обязан вручную строить backpressure через pause/resume и ограниченный буфер — или рискуешь OOM. Pull делает backpressure структурным и бесплатным — а Node Readable асинхронно-итерируем, так что for await (const chunk of readable) придушивает реальный I/O автоматически. Клади очистку ресурсов в finally: выход из цикла через break/return/throw вызывает return() итератора, который пропускает код после текущего yield, но выполняет finally, закрывая курсоры и освобождая соединения; помести эту очистку после цикла вместо этого — и первый же break утечёт соединение, пока пул не умрёт. Сцепляй генераторы (source → map → filter → consumer) ради ленивого конвейера, держащего по одному элементу на стадию — память O(1) при стриминге 10М строк вместо материализации ~1.5 ГБ массива — с pull-backpressure, прошитым из конца в конец. Только помни: for await последователен, а не конкурентен: если нужна параллельность, ограничь её пулом воркеров, а не неограниченным Promise.all. Теперь, когда видишь запрос, буферизующий все строки перед обработкой, знаешь фикс: заменить массив результатов на асинхронный генератор над курсором — та же логика, постоянная память и backpressure бесплатно.
Практика
Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.
Что-то непонятно?
Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.