Урок 09

Асинхронные итераторы

Обычные итераторы работают с синхронными данными. Асинхронные — с потоками данных, которые приходят со временем.

Проблема

// Не работает
const res = await fetch('/api/large');
const reader = res.body.getReader();

for (const chunk of reader) { // ошибка
    console.log(chunk);
}

reader — не итерируемый. Нужен асинхронный for await...of.

for await...of

const res = await fetch('/api/large');

for await (const chunk of res.body) {
    console.log(chunk);
}

res.body — ReadableStream, поддерживает асинхронную итерацию.

Async generator

async function* numbers() {
    for (let i = 1; i <= 3; i++) {
        await new Promise((r) => setTimeout(r, 1000));
        yield i;
    }
}

for await (const n of numbers()) {
    console.log(n); // 1, 2, 3 с задержкой
}

function* с async — асинхронный генератор.

Разница с обычным генератором

// Синхронный
function* sync() {
    yield 1;
    yield 2;
}

// Асинхронный
async function* asyncGen() {
    yield 1;
    yield 2;
}

Асинхронный возвращает промисы, await работает внутри.

Практика: пагинация

async function* fetchAllPages(baseUrl) {
    let page = 1;
    let hasMore = true;

    while (hasMore) {
        const res = await fetch(`${baseUrl}?page=${page}`);
        const data = await res.json();

        yield data.items;

        hasMore = data.hasMore;
        page++;
    }
}

for await (const items of fetchAllPages('/api/users')) {
    console.log('Страница:', items);
}

Итерируешь страницы, не думая о пагинации.

Практика: чтение файла

import { createReadStream } from 'fs';

async function* readLines(path) {
    const stream = createReadStream(path, { encoding: 'utf-8' });
    let buffer = '';

    for await (const chunk of stream) {
        buffer += chunk;
        const lines = buffer.split('\n');
        buffer = lines.pop() ?? '';

        for (const line of lines) {
            yield line;
        }
    }

    if (buffer) yield buffer;
}

for await (const line of readLines('large.log')) {
    console.log(line);
}

Прерывание

for await (const item of stream()) {
    if (item.done) break;
}

break останавливает итерацию. Стрим закрывается.

Практика: retry generator

async function* retryGenerator(fn, attempts = 3) {
    for (let i = 0; i < attempts; i++) {
        try {
            yield await fn();
            return;
        } catch (error) {
            if (i === attempts - 1) throw error;
            await new Promise((r) => setTimeout(r, 1000 * (i + 1)));
        }
    }
}

for await (const result of retryGenerator(() => fetchData())) {
    console.log(result);
}

Преобразование стрима

async function* transform(source, fn) {
    for await (const item of source) {
        yield fn(item);
    }
}

async function* filter(source, predicate) {
    for await (const item of source) {
        if (predicate(item)) yield item;
    }
}

const source = numbers();
const doubled = transform(source, (n) => n * 2);
const evens = filter(doubled, (n) => n % 4 === 0);

for await (const n of evens) {
    console.log(n);
}

Итерация SSE

async function* sse(url) {
    const res = await fetch(url);
    const reader = res.body.getReader();
    const decoder = new TextDecoder();
    let buffer = '';

    while (true) {
        const { done, value } = await reader.read();
        if (done) break;

        buffer += decoder.decode(value, { stream: true });
        const lines = buffer.split('\n\n');
        buffer = lines.pop() ?? '';

        for (const line of lines) {
            if (line.startsWith('data: ')) {
                yield JSON.parse(line.slice(6));
            }
        }
    }
}

for await (const event of sse('/api/stream')) {
    console.log(event);
}

Возврат и finally

async function* gen() {
    try {
        yield 1;
        yield 2;
    } finally {
        console.log('очистка');
    }
}

for await (const x of gen()) {
    if (x === 1) break;
}
// выведет 'очистка'

finally вызывается при break или return.

Итоги

  • for await...of для асинхронных итераторов
  • async function* для асинхронных генераторов
  • Пагинация через генератор
  • Чтение стримов построчно
  • Преобразование потоков
  • finally при прерывании