Раздел 25 · Effect-TS

Sink и продвинутый Stream: приёмники, пакеты, DLQ, состояние по ходу

senior~80 мин

открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти

Sink и продвинутый Stream: приёмники, пакеты, DLQ, состояние по ходу

Сцена · поток собрали, а складывать некуда

В 09 · Stream мы научились получать данные: тянуть из push-источника, буферизовать, окнить, мержить, раздавать многим потребителям. Логичный следующий вопрос: а куда всё это девать.

В жизни ответ звучит примерно так: писать в базу пачками по сто штук, но не реже раза в секунду; если база лежит, повторить три раза с растущей паузой; если и это не помогло, отложить в сторонку и не потерять; параллельно считать статистику для дашборда; а ещё уметь дописывать в файл.

Наивная реализация этого списка выглядит как for await (const item of stream) с четырьмя счётчиками, ручным буфером, таймером и try/catch внутри. Работает ровно до первого граничного случая: пачка недобрана, а поток кончился; таймер сработал во время записи; ошибка на третьем элементе, а первые два уже ушли.

Sink это вторая половина стрима: описание того, куда поток вливается. Такое же композируемое значение, как сам Stream.

Карта урока · что заберёшь домой

Восемь разделов:

  1. Раздел 1, что такое Sink и какие есть готовые.
  2. Раздел 2, свой приёмник через reduce, fold и forEach.
  3. Раздел 3, комбинирование приёмников: два результата за проход, orElse и его ловушка.
  4. Раздел 4, запись в файл и в базу, ретраи на уровне пачки.
  5. Раздел 5, aggregateWithin, пачка по размеру и по времени.
  6. Раздел 6, очередь недоставленных сообщений.
  7. Раздел 7, состояние по ходу потока: scan и mapAccum.
  8. Раздел 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 встречались тут постоянно. Следующий урок разбирает их и остальную стандартную библиотеку систематически.

Полезно перечитать: