Sink и продвинутый Stream: приёмники, пакеты, DLQ, состояние по ходу
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Sink и продвинутый Stream: приёмники, пакеты, DLQ, состояние по ходу
Сцена · поток собрали, а складывать некуда
В 09 · Stream мы научились получать данные: тянуть из push-источника, буферизовать, окнить, мержить, раздавать многим потребителям. Логичный следующий вопрос: а куда всё это девать.
В жизни ответ звучит примерно так: писать в базу пачками по сто штук, но не реже раза в секунду; если база лежит, повторить три раза с растущей паузой; если и это не помогло, отложить в сторонку и не потерять; параллельно считать статистику для дашборда; а ещё уметь дописывать в файл.
Наивная реализация этого списка выглядит как for await (const item of stream) с четырьмя счётчиками, ручным буфером, таймером и try/catch внутри. Работает ровно до первого граничного случая: пачка недобрана, а поток кончился; таймер сработал во время записи; ошибка на третьем элементе, а первые два уже ушли.
Sink это вторая половина стрима: описание того, куда поток вливается. Такое же композируемое значение, как сам Stream.
Карта урока · что заберёшь домой
Восемь разделов:
- Раздел 1, что такое
Sinkи какие есть готовые. - Раздел 2, свой приёмник через
reduce,foldиforEach. - Раздел 3, комбинирование приёмников: два результата за проход,
orElseи его ловушка. - Раздел 4, запись в файл и в базу, ретраи на уровне пачки.
- Раздел 5,
aggregateWithin, пачка по размеру и по времени. - Раздел 6, очередь недоставленных сообщений.
- Раздел 7, состояние по ходу потока:
scanиmapAccum. - Раздел 8, постраничный источник через
Stream.paginate.
Раздел 1 · Sink это тоже значение
Sink<A, In, L, E, R> читается по каналам почти как Effect: In это то, что приходит на вход, A это то, что получится в итоге, дальше остаток, ошибка и требования.
Готовых приёмников хватает на большинство задач:
import { Effect, Sink, Stream } from 'effect';
const total = yield* events.pipe(Stream.run(Sink.count));
const sum = yield* events.pipe(Stream.map((e) => e.status), Stream.run(Sink.sum));
const first3 = yield* events.pipe(Stream.run(Sink.take(3)));
const all = yield* events.pipe(Stream.run(Sink.collect()));
// count = 10 | sum = 2600 | first3 = 3 элемента | all = 10 элементов
Полезные из коробки: collect, take, count, sum, head, last, find, drain, forEach, fold, reduce, takeWhile, takeUntil, timed, never.
Заметь две мелочи в именах, на которых легко споткнуться. Sink.collect() это функция, а не значение, поэтому скобки обязательны. А набрать ровно N элементов теперь Sink.take(n), симметрично Stream.take, вместо старого collectAllN.
Два наблюдения, ради которых стоило вводить отдельное понятие.
Sink.take(3) останавливает поток, когда набрал три элемента. Приёмник управляет чтением, а не просто пассивно принимает: это та же pull-модель, только с другой стороны трубы.
Sink.count не собирает элементы в память. Пробежать миллиард элементов и получить число это O(1) по памяти. Sink.collect() наоборот, поэтому на неограниченном потоке он вешает процесс.
Все приёмники, собирающие много элементов, отдают обычный Array. Chunk из пользовательского API потоков в v4 практически исчез: он остался внутренней деталью, а наружу торчат массивы.
Раздел 2 · Свой приёмник
Два способа, и выбор между ними механический.
Sink.reduce, когда нужен результат. Сворачиваем поток в одно значение за один проход:
const errorRate = Sink.reduce(
() => ({ ok: 0, failed: 0 }), // начальное состояние
(acc, event: Event) =>
event.status >= 500 ? { ...acc, failed: acc.failed + 1 } : { ...acc, ok: acc.ok + 1 },
);
const rate = yield* events.pipe(Stream.run(errorRate));
// { ok: 8, failed: 2 }
Начальное состояние это функция, а не значение: иначе один изменяемый объект утёк бы между прогонами.
Для ранней остановки есть Sink.reduceWhile, у которого вторым аргументом идёт предикат “продолжаем ли”. Условие (acc) => acc.failed < 10 перестанет читать, как только накопится десять ошибок, и остаток потока даже не будет запрошен.
Если шаг свёртки сам должен быть эффектом (сходить в базу, посмотреть на часы), берут Sink.fold с той же формой: он отличается от reduce ровно тем, что функция шага возвращает Effect.
Sink.forEach, когда нужен побочный эффект. Записать, отправить, положить в очередь:
const toDatabase = Sink.forEach((event: Event) => repository.insert(event));
Ещё один комбинатор, который экономит много кода, это Sink.mapInput: он адаптирует вход приёмника. Написал один Sink, который умеет писать строки, и переиспользуй его для любого потока, добавив преобразование:
const lines = Sink.forEach((line: string) => appendToFile(line));
const jsonl = Sink.mapInput(lines, (event: Event) => `${JSON.stringify(event)}\n`);
Это тот же приём, что contramap в функциональных библиотеках: Stream.map меняет данные до приёмника, Sink.mapInput меняет то, что приёмник согласен принимать. Разница чисто в том, где живёт знание о формате: в конвейере или в приёмнике.
Раздел 3 · Комбинирование приёмников
Два результата за один проход
Считать две величины хочется, а прогонять поток дважды нельзя: источник бывает одноразовым. Готового Sink.zip в v4 нет, и не нужен: две величины это просто одно составное состояние.
const both = Sink.reduce(
() => ({ n: 0, ok: 0, failed: 0 }),
(acc, event: Event) => ({
n: acc.n + 1,
ok: acc.ok + (event.status < 500 ? 1 : 0),
failed: acc.failed + (event.status >= 500 ? 1 : 0),
}),
);
const stats = yield* events.pipe(Stream.run(both));
// { n: 10, ok: 8, failed: 2 }
Поток прочитан один раз, величин три. Приём общий: если приёмники это чистые свёртки, склеивай их состояния в один объект. Если же нужно именно два независимых потребителя с побочными эффектами (записать в базу и отправить в очередь), поток разветвляют явно через Stream.broadcast, и это честнее: видно, что потребителей действительно двое.
Гонка приёмников
Отдельного Sink.race тоже больше нет. Задача “верни первые сто записей или всё, что накопилось за пять секунд, что наступит раньше” решается на уровне потока: Stream.groupedWithin(100, '5 seconds') или, если условие полной пачки сложнее, Stream.aggregateWithin из Раздела 5. Это и читается яснее: ограничение по времени принадлежит трубе, а не приёмнику.
Sink.orElse и его ловушка
Кажется, что запасной приёмник это Sink.orElse(основной, () => запасной). Проверим на потоке из шести элементов, где основной падает начиная с третьего:
yield* events.pipe(Stream.rechunk(1), Stream.run(Sink.orElse(primary, () => fallback)));
// запасной увидел: [4, 5, 6]
Тройка потеряна. Основной приёмник её уже прочитал, упал на ней, и запасному досталось только то, что осталось непрочитанным. А если убрать Stream.rechunk(1) и оставить один большой чанк, запасной не увидит вообще ничего: весь чанк ушёл в основной приёмник за один заход.
Sink.orElse это про “если этот приёмник не справился, дочитай остаток другим”, а не про “перезалей то же самое в другое место”.
Настоящий запасной путь для доставки делается на уровне эффекта, внутри forEach:
const write = (event: Event) =>
primaryStore.insert(event).pipe(
Effect.orElse(() => fallbackStore.insert(event)),
);
yield* events.pipe(Stream.run(Sink.forEach(write)));
// основной: [1, 2] | запасной: [3, 4, 5, 6]
Теперь ни один элемент не потерялся: упавший на основном хранилище уходит в запасное. Правило то же, что мы вывели в 22 · Schema как контракт для модели: решение о запасном пути принимается на уровне элемента, а не на уровне трубы.
Раздел 4 · Файл, база, ретраи
Файл
FileSystem из 20 · Config, секреты и platform отдаёт готовый приёмник:
import { FileSystem, Stream } from 'effect';
const fs = yield* FileSystem.FileSystem;
yield* events.pipe(
Stream.map((event) => `${JSON.stringify(event)}\n`),
Stream.encodeText,
Stream.run(fs.sink(file)),
);
fs.sink(path) принимает Uint8Array, поэтому перед ним ставится Stream.encodeText. Файл открывается один раз на весь поток и закрывается при завершении, включая аварийное: приёмник живёт в Scope, как всё из 05 · Resources.
База и ретраи на уровне пачки
Вставлять по одной записи в базу это способ убить и базу, и латентность. Пакетная вставка с ретраями:
const insertBatch = (batch: ReadonlyArray<Event>) =>
repository.insertMany(batch).pipe(
Effect.retry({ times: 5, schedule: Schedule.exponential('50 millis') }),
Effect.as(batch.length),
);
const inserted = yield* events.pipe(
Stream.grouped(100),
Stream.mapEffect(insertBatch),
Stream.runFold(0, (acc, n) => acc + n),
);
Stream.grouped(100) отдаёт непустой ReadonlyArray, а не Chunk, поэтому .length и обычные методы массива работают напрямую, без переходников вроде Chunk.toReadonlyArray.
Ключевое слово тут на уровне пачки. Ретрай стоит именно на insertMany, а не на всём потоке: упавшая вставка повторяется, уже вставленные пачки не переигрываются. Ретрай на уровне потока переиграл бы всё с начала, а это дубликаты.
Из этого следует требование к самой вставке: она должна быть идемпотентной. Сбой мог произойти после записи, но до ответа, и повтор запишет то же самое второй раз.
Раздел 5 · Пачка по размеру и по времени
Stream.grouped(100) из предыдущего раздела ждёт ровно сто элементов. Если поток притих на девяноста девяти, они зависнут до конца времён. Правильное поведение: “сто штук или раз в секунду, что наступит раньше”.
Именно это делает Stream.aggregateWithin, и приёмник в нём это первый аргумент:
const batches = yield* events.pipe(
Stream.aggregateWithin(Sink.take<Event>(4), Schedule.spaced('50 millis')),
Stream.map((batch) => batch.length),
Stream.runCollect,
);
// размеры пачек: [4, 4, 2]
Последняя пачка неполная, и это правильно: поток кончился, ждать больше нечего.
Как читать сигнатуру: Sink определяет, что считается полной пачкой, Schedule определяет, когда сдаться и отдать неполную. Оба параметра любые, так что “пачка не больше мегабайта суммарного веса” делается через Sink.reduceWhile с весом в состоянии, а не костылями поверх счётчика.
С этим же комбинатором в 09 · Stream мы встречались в виде groupedWithin. Разница в том, что groupedWithin умеет только “N штук за T”, а aggregateWithin принимает произвольный приёмник и произвольное расписание.
Раздел 6 · Очередь недоставленных
Элемент, который не удалось обработать, нельзя ни потерять, ни бесконечно повторять. Стандартный ответ это dead-letter queue: отложить в сторону, пометить причиной, идти дальше.
const dlq = yield* Queue.unbounded<Event>();
const delivered = yield* events.pipe(
Stream.mapEffect((event) =>
process(event).pipe(
Effect.as(Option.some(event)),
Effect.catch((error) =>
Effect.logWarning(`отправляю в DLQ: ${error.message}`).pipe(
Effect.zipRight(Queue.offer(dlq, event)),
Effect.as(Option.none<Event>()),
),
),
),
),
Stream.filterMap((option) => option),
Stream.runCollect,
);
// доставлено = 7 | в DLQ = 2
Три вещи, которые стоит сделать сразу, а не потом.
Класть в очередь не голое событие, а событие плюс причину плюс число попыток. Без этого разбор DLQ превращается в гадание.
Ограничивать очередь. Queue.unbounded хорош в примере и плох в проде: если сломалось всё, неограниченная очередь съест память. Queue.dropping(10_000) плюс метрика длины из 18 · Observability честнее.
Предусмотреть переигрывание. DLQ без кнопки “прогнать заново после починки” это просто аккуратная свалка.
Раздел 7 · Состояние по ходу потока
Иногда решение зависит от истории. Не от всего потока (это fold и конец), а от того, что было до текущего элемента.
const running = yield* events.pipe(
Stream.mapAccum(0, (failures, event) => {
const next = event.status >= 500 ? failures + 1 : 0;
return [next, { id: event.id, consecutiveFailures: next }] as const;
}),
Stream.filter((row) => row.consecutiveFailures > 0),
Stream.runCollect,
);
Stream.mapAccum(seed, step) тащит состояние вдоль потока: step получает текущее состояние и элемент, возвращает новое состояние и то, что уйдёт дальше. Здесь так считается счётчик подряд идущих провалов, то самое, на чём срабатывает circuit breaker.
Соседи по семейству:
Stream.scan(seed, f)отдаёт каждое промежуточное состояние. Бегущая сумма, накопительный график;Stream.mapAccumEffect, то же, чтоmapAccum, но шаг может быть эффектом (сходить в базу, посмотреть на часы);Stream.zipWithPreviousдаёт пару “предыдущий, текущий”, когда истории нужно ровно на один шаг.
Главное отличие от Ref снаружи потока: состояние здесь локально для конкретного прогона. Два параллельных прогона не мешают друг другу, и очищать ничего не надо.
Раздел 8 · Постраничный источник
Симметрично приёмникам, у потока бывает нетривиальный источник. Самый частый в жизни это API с курсором:
const all = yield* Stream.paginate(1, (page) =>
api.fetchPage(page).pipe(Effect.map(({ items, nextPage }) => [items, nextPage] as const)),
).pipe(Stream.runCollect);
// со всех страниц: [1, 2, 3, 4, 5]
Stream.paginate(seed, step) вызывает step, пока тот возвращает Option.some(следующий курсор). Обрати внимание на форму шага в v4: он возвращает эффект с парой “массив элементов страницы, следующий курсор”. То есть страница целиком отдаётся за один шаг, и разворачивать её отдельным flattenIterable больше не надо, поток сам раскладывает массив на элементы.
Что тут выигрывается по сравнению с циклом while. Backpressure: следующая страница запрашивается, только когда потребитель разобрался с текущей. Композиция: сверху навешивается RateLimiter из 19 · HttpClient и обвязка API, ретраи, таймаут, метрики. И остановка: Option.none() заканчивает поток, никаких флагов и break.
Pulse · вклад этого урока
Запись событий переезжает с самописного буфера на конвейер:
// pulse-<nick>/src/pipeline.ts
import { Effect, FileSystem, Metric, Queue, Schedule, Sink, Stream } from 'effect';
export const runPipeline = Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const storage = yield* Storage;
const config = yield* PulseConfig;
const dlq = yield* Queue.dropping<FailedEvent>(10_000);
const persist = (batch: ReadonlyArray<ProbeEvent>) =>
storage.insertMany(batch).pipe(
Effect.retry({ times: 3, schedule: Schedule.exponential('100 millis') }),
Effect.catch((error) =>
Queue.offer(dlq, { batch, reason: error.message }).pipe(
Effect.zipRight(Effect.logError('пачка не записалась, ушла в DLQ')),
),
),
Effect.as(batch.length),
);
yield* monitorEventsStream.pipe(
Stream.aggregateWithin(Sink.take<ProbeEvent>(100), Schedule.spaced(config.flushInterval)),
Stream.tap((batch) => Metric.update(eventsWritten, batch.length)),
Stream.mapEffect(persist),
Stream.runDrain,
);
});
Что это дало по сравнению с прошлой версией.
Буфер, таймер и счётчик исчезли: их роль играет aggregateWithin. Неполная пачка на завершении больше не теряется. Ретрай стоит на пачке, а не на всём потоке, поэтому повтор не создаёт дубликатов. Провалившаяся пачка уходит в ограниченную очередь вместе с причиной, а не в лог с надписью “не смогли”.
Плюс два вспомогательных приёмника: одна свёртка Sink.reduce считает статистику окна (сколько событий, сколько ошибок, максимальная латентность) тем же проходом, без второго обхода потока, а fs.sink пишет сырой JSONL параллельно с записью в базу, как страховка на время миграции хранилища.
Финал · чек-лист
ДЗ
Все задания делаются в pulse-<nick>. После каждого открой PR с тегом lesson-23.
Дальше
- 24 · Стандартная библиотека Effect.
Option,DurationиArrayвстречались тут постоянно. Следующий урок разбирает их и остальную стандартную библиотеку систематически.
Полезно перечитать:
- 09 · Stream, первая половина трубы. Теперь видно, почему источник и приёмник это симметричные понятия.
- 26-data-engineering · 07 Потоковая обработка, те же идеи (окна, backpressure, DLQ) в масштабе кластера, а не одного процесса.