Stream: pull vs push, asyncPush, buffer, grouped, merge, broadcast
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Stream: pull vs push, asyncPush, buffer, grouped, merge, broadcast
Сцена · WebSocket, который тебя топит
Ты пишешь сервис, который слушает WebSocket с событиями (клики, тики биржи, что угодно). На каждое событие нужно сходить во внешний API. И вот ты пишешь самое естественное, что приходит в голову:
socket.on('message', async (event) => {
await processEvent(event);
});
Через секунду тебя топит. Источник присылает 100 событий в секунду, processEvent стоит 500 мс на запрос, и теперь в полёте 500 параллельных HTTP-запросов. Внешний API возвращает 429. Память течёт. На графиках красное.
Корень проблемы в том, что producer решает, когда отдавать данные, а у consumer нет рычага сказать “погоди”. На любом push-источнике (WebSocket, EventEmitter, setInterval) ты получаешь ровно три плохих варианта:
- дропать события (теряешь данные);
- копить в очереди (память растёт неограниченно);
- запускать параллельно (систему разносит).
Сравни с пагинацией:
const page1 = await repo.getPage(1);
// обрабатываем page1
const page2 = await repo.getPage(2);
База не вываливает миллион записей. Она ждёт, пока ты попросишь следующую страницу. Темп задаёт consumer. Это и есть pull-модель, и Effect Stream это её реализация на тип-системе.
В этом уроке разбираем, как Stream<A, E, R> устроен внутри, как затянуть push-источник в pull, как буферизовать, окнить, мержить, и как раздать один стрим многим consumer-ам.
Карта урока · что заберёшь домой
Восемь разделов:
- Раздел 1, pull-based stream и почему backpressure получается бесплатно.
- Раздел 2,
mapEffectсconcurrency, контролируемый параллелизм без потери pull-семантики. - Раздел 3, cold vs hot, что значит “Stream ленив” и почему
Stream.iterate(0, n => n + 1)не вешает рантайм. - Раздел 4,
Stream.callback, мост между push-источниками и pull-стримом.bufferSize,dropping,sliding. - Раздел 5,
Stream.buffer,Stream.grouped,Stream.groupedWithin. Окна по числу и по времени. - Раздел 6,
Stream.merge,Stream.mergeAll. Halt-стратегии,concurrency,bufferSize. - Раздел 7,
broadcastN,broadcast,share. Один источник, много consumer-ов, replay и idle-cleanup. - Раздел 8, Pulse:
monitorEventsStream,groupedWithinдля пакетной записи JSONL,shareдля UI и подписчиков.
К концу урока ты пишешь стрим из любого источника, отличаешь backpressure от дроп-стратегии, выбираешь между broadcast и share за 30 секунд, и понимаешь, почему Stream.buffer после Stream.callback обычно бесполезен.
Раздел 1 · Pull-based stream и backpressure из ничего
Источник, который ждёт consumer-а
Возьмём пагинацию из сцены и сделаем из неё настоящий стрим:
import { Console, Effect, Option, Stream } from 'effect';
type User = { readonly id: number; readonly name: string };
const db: Record<number, { users: ReadonlyArray<User>; nextPage: Option.Option<number> }> = {
1: { users: [{ id: 1, name: 'Alice' }, { id: 2, name: 'Bob' }], nextPage: Option.some(2) },
2: { users: [{ id: 3, name: 'Charlie' }, { id: 4, name: 'Diana' }], nextPage: Option.some(3) },
3: { users: [{ id: 5, name: 'Eve' }], nextPage: Option.none() },
};
const fetchPage = (page: number) =>
Console.log(`[DB] fetch page ${page}`).pipe(Effect.as(db[page]!));
const users = Stream.paginate(1, (page) =>
fetchPage(page).pipe(Effect.map(({ users, nextPage }) => [users, nextPage] as const)),
);
// users: Stream<User, never, never>
Stream.paginate(seed, step) повторно вызывает step(seed), ожидая Effect<[ReadonlyArray<A>, Option<nextSeed>]>. Пока Option.some, продолжает; на Option.none останавливается. Массив из шага разворачивается в отдельные элементы стрима сам, отдельный flatten не нужен.
Обрати внимание на форму шага: он и эффектный, и пачечный сразу. Раньше в Effect для этого было два разных конструктора (paginate для чистого шага по одному элементу и paginateEffect для эффектного), теперь остался один, который покрывает оба случая. Нужен один элемент за шаг, верни массив из одного элемента.
Теперь дёрнем его медленным consumer-ом:
const program = users.pipe(
Stream.tap((user) => Console.log(`[USE] ${user.name}`)),
Stream.tap(() => Effect.sleep('100 millis')),
Stream.runDrain,
);
await Effect.runPromise(program);
// [DB] fetch page 1
// [USE] Alice
// [USE] Bob
// [DB] fetch page 2
// [USE] Charlie
// [USE] Diana
// [DB] fetch page 3
// [USE] Eve
Запомни порядок: страница 1 fetched, два пользователя обработаны, только потом fetched страница 2. Источник не дёргает следующую страницу, пока consumer не разобрался с текущей.
Это и есть pull. Никаких специальных “stop”-сигналов, никаких pause/resume API. Просто рантайм спрашивает у источника “дай элемент” каждый раз, когда downstream его попросил. Если downstream спит, источник тоже спит.
Backpressure это не фича, это форма
В push-системах backpressure это отдельный механизм: producer должен уметь приостанавливаться, consumer должен уметь сигнализировать. Стандарт WHATWG Streams (05 · Async/06) ввёл это как часть протокола. Node Streams используют highWaterMark и события drain.
В pull-модели backpressure это отсутствие приглашения. Источник не может производить быстрее, чем consumer pull-ит. Никакого протокола не нужно, это просто следствие “consumer задаёт темп”.
Когда тебе говорят “Effect Stream имеет встроенный backpressure”, это не маркетинг, это буквально определение pull-стрима.
Что под капотом
Stream<A, E, R> это, грубо говоря, Effect<NonEmptyReadonlyArray<A>, E, R>, обёрнутый в скоуп. Такой эффект в Effect называется Pull, и увидеть его можно через Stream.toPull. На каждый pull рантайм запускает этот эффект:
- вернул непустой массив, значит вот тебе порция элементов, продолжай дёргать;
- упал со специальным сигналом
Halt, значит источник иссяк, стрим закончен; - упал с ошибкой
E, стрим завершается с этой ошибкой.
Порционность важна: тянуть пачками от 16 до 256 элементов производительнее, чем по одному. Поэтому большинство трансформаций (map, filter) работают над пачкой целиком. Пачка это обычный массив: отдельный тип Chunk в Effect ещё есть, но из публичного Stream-API он ушёл, наружу везде выходят массивы.
Технически это всё построено на Channel, более низкоуровневой абстракции. Channel поверх (входной chunk типа In, выходной Out, окончание Done) умеет описывать любой стрим, sink, и пайплайн между ними. Для большинства задач Channel напрямую не нужен, но он есть, если понадобится свой кастомный примитив.
Что взять с собой
Stream<A, E, R>это ленивая последовательность Effect-ов, в pull-модели consumer задаёт темп.Stream.paginateэто идиоматический способ затянуть страничный API в стрим: шаг эффектный и отдаёт сразу пачку.- Backpressure встроен в природу pull: пока downstream спит, источник тоже спит.
- Под капотом порции-массивы и
Channel, но в 95 процентах задач достаточноStream-API.
Раздел 2 · mapEffect с concurrency
Последовательно это часто медленно
В прошлом разделе Stream.tap(() => Effect.sleep('100 millis')) гарантировал, что элементы обрабатываются строго один за другим. Это удобно для отладки, но в реальной задаче “сходить во внешний API на каждого пользователя” последовательность это потеря времени.
Stream.mapEffect запускает Effect-функцию на каждом элементе. С опцией concurrency ты задаёшь, сколько Effect-ов рантайм может держать в полёте одновременно:
import { Console, Effect, Option, Stream } from 'effect';
type User = { readonly id: number; readonly name: string };
const fetchPage = (page: number) =>
Console.log(`[DB] fetch page ${page}`).pipe(Effect.as(db[page]!));
const users = Stream.paginate(1, (page) =>
fetchPage(page).pipe(Effect.map(({ users, nextPage }) => [users, nextPage] as const)),
);
const processUser = (user: User) =>
Console.log(`[PROC] start ${user.name}`).pipe(
Effect.andThen(Effect.sleep(`${100 + Math.random() * 100} millis`)),
Effect.andThen(Console.log(`[PROC] done ${user.name}`)),
Effect.as(user),
);
const program = users.pipe(
Stream.mapEffect(processUser, { concurrency: 3 }),
Stream.runDrain,
);
В логе увидишь:
[DB] fetch page 1
[PROC] start Alice
[PROC] start Bob
[DB] fetch page 2
[PROC] start Charlie
[PROC] done Alice
[PROC] start Diana
[PROC] done Bob
[DB] fetch page 3
[PROC] start Eve
Три обработчика крутятся параллельно, и бэкпрешер всё ещё работает: страница 2 подтягивается ровно тогда, когда стриму нужно подкормить свободные слоты. Внутри mapEffect сидит семафор на N пермитов: пока все заняты, новые элементы из upstream не запрашиваются.
Сохранение порядка
По умолчанию mapEffect({ concurrency: N }) сохраняет порядок выхода: даже если Diana обработалась раньше Charlie, в downstream они придут в исходном порядке. Это удобно для записи в файл/БД “в том же порядке, что пришло”.
Если порядок не важен и хочется выжать максимум, есть unordered: true:
Stream.mapEffect(processUser, { concurrency: 8, unordered: true });
Тогда элемент уходит downstream сразу, как только его Effect завершился. Чуть быстрее, чуть прожорливее по памяти на буфер переупорядочивания (которого теперь нет).
Что взять с собой
Stream.mapEffect(fn, { concurrency })запускает доNEffect-ов параллельно, оставаясь pull-стримом.- Backpressure через семафор внутри: новые элементы тянутся только если есть свободный пермит.
- По умолчанию порядок выхода равен порядку входа;
unordered: trueотпускает это ограничение.
Раздел 3 · Cold stream, ленивая декларация, бесконечные последовательности
Hot vs Cold
setInterval это hot источник: ты вызвал, и оно тикает. Хочешь остановить, вспоминай clearInterval. Забыл, получил утечку.
Stream это cold источник. Описание стрима ничего не делает, это просто значение. Только Effect.runPromise(Stream.run*(...)) приводит описание в исполнение:
import { Effect, Stream } from 'effect';
const infinite = Stream.iterate(0, (n) => n + 1);
// ничего не происходит, infinite это просто данные
const finite = Stream.take(infinite, 5);
// тоже ничего не происходит, finite это уточнённое описание
const result = await Effect.runPromise(Stream.runCollect(finite));
// здесь рантайм впервые что-то делает: pull 5 элементов и стоп
console.log(result); // [0, 1, 2, 3, 4]
Stream.runCollect отдаёт обычный Array<A>. Держи это в голове, если читаешь старые примеры: раньше он возвращал Chunk, и вокруг каждого результата стоял Chunk.toReadonlyArray.
Stream.iterate(0, n => n + 1) это бесконечный стрим, но он безопасен, потому что cold: пока никто не запросил, никаких ресурсов он не ест. Stream.take(_, 5) ограничивает потребление пятью элементами, после которых рантайм закрывает стрим и освобождает скоуп.
Это та же дисциплина, что у Effect: значение Effect<A, E, R> это описание, runPromise это запуск. У Stream ровно та же двухфазность.
Любой стрим заворачивается в скоуп
Stream.run* всегда работает в скоупе. Если стрим открывает ресурс (файл, соединение, listener), он гарантированно закроется при завершении стрима, даже если consumer прервал работу через Fiber.interrupt:
const file = Stream.unwrap(
Effect.acquireRelease(
Effect.sync(() => openFile('events.jsonl')),
(fd) => Effect.sync(() => fd.close()),
).pipe(
Effect.map((fd) =>
Stream.fromReadableStream({
evaluate: () => fd.readable,
onError: (cause) => new ReadError({ cause }),
}),
),
),
);
Stream.unwrap берёт Effect, который отдаёт стрим, и приклеивает его скоуп к скоупу стрима: Scope из требований вычитается прямо тут. Отдельного unwrapScoped в Effect больше нет, обычный unwrap покрывает оба случая.
Это та же acquireRelease-дисциплина, что в 05 · Resources. При любом исходе (успех, ошибка, interrupt) fd.close() будет вызван.
Что взять с собой
Streamэто cold-источник: декларация не делает ничего, исполнение начинается только наrun*.- Бесконечные стримы безопасны, пока ты ограничиваешь потребление (
take,takeUntil,takeWhile). - Скоуп закрывается автоматически: файлы, соединения, listener-ы освобождаются по завершении стрима.
Раздел 4 · Stream.callback, мост из push в pull
Зачем мост
Не всё в мире pull. WebSocket, EventEmitter, setInterval, click-events, push-уведомления, очередь сообщений с auto-acknowledge: всё это push-источники. Чтобы их затянуть в Stream, нужен мост, который сам сообщает, когда у источника появилось новое событие.
Stream.callback это такой мост. Ты даёшь ему функцию подписки, получаешь обычный pull-стрим:
import { Effect, Queue, Stream } from 'effect';
const ticks = Stream.callback<number>((queue) =>
Effect.acquireRelease(
Effect.sync(() => {
let count = 0;
return setInterval(() => Queue.offerUnsafe(queue, count++), 100);
}),
(handle) => Effect.sync(() => clearInterval(handle)),
),
);
Что происходит:
Stream.callbackзовёт твою функцию в момент, когда первый consumer начинает pull-ить, и передаёт ей внутреннююQueue;- внутри
Effect.acquireReleaseты подписываешься на источник (setInterval) и регистрируешь cleanup (clearInterval); - каждый внешний event превращается в
Queue.offerUnsafe(queue, value), который кладёт значение в эту очередь; - стрим вытягивает значения из очереди по своему темпу.
Никакого специального emit-объекта тут нет, работаешь обычным модулем Queue из 07 · Координация:
Queue.offerUnsafe(queue, value), одно значение из синхронного колбэка;Queue.offerAllUnsafe(queue, [1, 2, 3]), сразу пачка;Queue.offer(queue, value), эффектная версия, которая честно ждёт место в буфере;Queue.endUnsafe(queue), стрим закончился;Queue.failCauseUnsafe(queue, Cause.fail(error)), стрим упал (эффектный вариант короче:Queue.fail(queue, error)).
Суффикс Unsafe тут означает ровно одно: вызов синхронный и ничего не ждёт, его можно позвать из колбэка любой сторонней библиотеки. Такая конвенция в Effect сквозная, Unsafe всегда суффикс.
Буфер по умолчанию неограничен
И вот тут ловушка. По умолчанию callback использует unbounded буфер. Если источник эмитит быстрее, чем consumer тянет, буфер растёт. Память течёт ровно так же, как в наивном socket.on('message', ...). Pull-семантика на самом нижнем слое спасает downstream, но внутренняя очередь callback это уже push-зона.
Решается опциями:
Stream.callback<number>(
(queue) => Effect.acquireRelease(
Effect.sync(() => {
let count = 0;
return setInterval(() => Queue.offerUnsafe(queue, count++), 100);
}),
(handle) => Effect.sync(() => clearInterval(handle)),
),
{ bufferSize: 4, strategy: 'dropping' },
);
Три стратегии переполнения:
'suspend'(по умолчанию для bounded), producer подвисает до освобождения места. ДляsetInterval-стиля бесполезно: timer-колбэк не умеет “подождать”, аQueue.offerUnsafeпо определению не ждёт и просто вернётfalse;'dropping', при заполненном буфере новые элементы (точнее, новые батчи) отбрасываются, downstream продолжает читать то, что уже было;'sliding', при заполненном буфере старые элементы выбрасываются, новые занимают их место. Полезно, когда важно “последнее значение”, не “первое”.
Тонкость · батчинг внутри буфера
Очередь под callback группирует подряд идущие offerUnsafe в один batch для эффективности. bufferSize это число batch-ей в очереди, не число отдельных элементов. Поэтому в логах часто видишь “странные” гэпы:
offer 0
Consumed 0
offer 1, offer 2, offer 3, offer 4
Consumed 1 ← consumer проснулся, забрал batch
offer 5, offer 6, offer 7, ...
Consumed 5 ← следующий batch начинается с 5, элементы 2..4 потерялись
С dropping теряются целые батчи, с sliding старые батчи выкидываются под новые. Это нюанс, который ловит только при тестах, не на дев-машине.
Второй сигнал того же рода: Queue.offerUnsafe возвращает boolean, а Queue.offerAllUnsafe возвращает массив непринятых элементов. Если тебе важно знать, сколько ты потерял, эти возвраты и есть твой счётчик потерь.
Что взять с собой
Stream.callback(fn)это мост из push-источника (WebSocket, EventEmitter, setInterval) в pull-стрим.- Внутрь колбэка приходит обычная
Queue, кладёшь черезQueue.offerUnsafe, закрываешь черезQueue.endUnsafe. Effect.acquireReleaseвнутри гарантирует cleanup при завершении стрима.- По умолчанию буфер unbounded, и это утечка-кандидат. Сразу задавай
{ bufferSize, strategy }. droppingдропает новые батчи,slidingсдвигает окно,suspendждёт (для timer-источника бесполезен).
Раздел 5 · Stream.buffer, Stream.grouped, Stream.groupedWithin
Pull даёт темп downstream-у, но иногда хочется сглаживать ритм между этапами. Источник тянет данные пачками по 50 мс, обработчик жуёт по 500 мс. Если идти строго в такт, источник простаивает 450 мс из 500. Хочется prefetch: пока обработчик жуёт, источник заранее тянет следующие элементы в буфер. И симметрично: если источник внезапно замолчал, хочется не отдавать в БД по одной строке, а копить пачку и писать большим блоком, по числу или по времени.
В этом разделе три инструмента: buffer для prefetch между этапами, grouped и groupedWithin для сборки пачек на выход.
Stream.buffer, прослойка между этапами
Stream.buffer({ capacity }) ставит между upstream и downstream отдельный файбер-producer, который тянет из upstream и кладёт в очередь. Downstream берёт из этой очереди. На pull-источнике это даёт “забегание вперёд”: producer заранее тянет элементы, чтобы consumer не ждал.
import { Effect, Schedule, Stream } from 'effect';
const program = Stream.range(1, 100).pipe(
Stream.buffer({ capacity: 4 }),
Stream.tap((n) => Effect.log(`Consumed ${n}`)),
Stream.schedule(Schedule.spaced('500 millis')),
Stream.take(10),
Stream.runDrain,
);
Stream.schedule(Schedule.spaced(d)) это идиоматический способ задать темп downstream. Альтернатива Stream.tap(() => Effect.sleep(d)) работает так же, но Schedule богаче: экспоненциальная задержка, jitter, лимиты на число итераций.
Стратегии у buffer те же три, что у callback:
'suspend'(по умолчанию), файбер-producer висит на полной очереди, ждёт consumer;'dropping', новые элементы дропаются;'sliding', старые элементы вытесняются.
Подвох · Stream.buffer после callback бесполезен
Очень частая ошибка. Ты делаешь:
const ticks = Stream.callback<number>((queue) => /* setInterval */);
const program = ticks.pipe(
Stream.buffer({ capacity: 4, strategy: 'dropping' }), // ← вот это не работает
// ... медленный consumer
);
Логика автора: “ограничу буфер на этапе, защищусь от утечки”. Реальность: к моменту, когда Stream.buffer начинает тянуть из upstream, callback уже всё собрал в свою внутреннюю очередь, которая по умолчанию unbounded. Stream.buffer просто перекачивает из одной очереди в другую, ничего не блокируя у источника.
Правило: если источник push, лимит ставь на callback ({ bufferSize, strategy }), не на отдельный Stream.buffer после.
Stream.buffer хорош для pull-источников и для разнесения “медленный producer плюс быстрый consumer” в разные файберы (типа prefetch).
Stream.grouped(n), батч по числу
import { Effect, Stream } from 'effect';
const grouped = Stream.range(1, 25).pipe(Stream.grouped(5));
// grouped: Stream<NonEmptyReadonlyArray<number>>
await Effect.runPromise(
grouped.pipe(
Stream.tap((batch) => Effect.log(`batch: ${batch.join(',')}`)),
Stream.runDrain,
),
);
// batch: 1,2,3,4,5
// batch: 6,7,8,9,10
// ...
Простая группировка по n. Если в источнике остался “хвост” меньше n, он улетит в downstream целиком (не отбросится). Обрати внимание на тип элемента: это обычный массив, причём непустой на уровне типа. Никакого Chunk.toReadonlyArray вокруг батча писать не надо, все методы массива доступны сразу.
Stream.groupedWithin(n, duration), батч по числу или по времени
В реальности batching по числу одного мало. Если события приходят медленно (10 в час), ждать n=100 чтобы записать в БД, это ждать неделю. Хочется “100 штук или раз в 5 секунд, что наступит раньше”.
import { Effect, Queue, Stream } from 'effect';
const events = Stream.callback<number>(
(queue) =>
Effect.acquireRelease(
Effect.sync(() => {
let count = 0;
return setInterval(() => Queue.offerUnsafe(queue, count++), 50 + Math.random() * 150);
}),
(handle) => Effect.sync(() => clearInterval(handle)),
),
{ bufferSize: 64, strategy: 'sliding' },
);
const program = events.pipe(
Stream.groupedWithin(10, '1 second'),
Stream.take(5),
Stream.tap((batch) => Effect.log(`batch of ${batch.length}`)),
Stream.runDrain,
);
groupedWithin(10, '1 second') отдаёт батч в downstream, как только выполнено одно из:
- в нём накопилось 10 элементов;
- прошла 1 секунда с момента, как пришёл первый элемент батча.
Это идеальный примитив для:
- батчинга
INSERTв БД; - агрегации метрик (“каждую секунду или каждые 1000 событий”);
- rate-limit-а вызовов внешнего API.
В Pulse этот примитив будет ровно тем, который собирает события мониторинга в JSONL-блоки.
Что взять с собой
Stream.bufferэто отдельный файбер для prefetch между этапами, ставит свою очередь и три стратегии переполнения.- После
Stream.callbackотдельныйStream.bufferобычно бесполезен, лимит ставь наcallback. Stream.grouped(n)это батч строго по числу, “хвост” уходит без потерь.Stream.groupedWithin(n, d)это батч по числу или по времени, основной инструмент для async-агрегации.- Оба отдают downstream обычный массив, а не
Chunk.
Раздел 6 · merge и mergeAll
Два независимых стрима в один
Чат-бот с LLM. Один поток это токены ответа, второй это heartbeat-ping каждые несколько секунд (чтобы клиент не закрыл соединение). Они независимы по времени, и оба нужно отправлять клиенту через один канал.
Stream.merge(s1, s2) запускает оба стрима параллельно и выпускает элементы downstream по мере появления:
import { Effect, Schedule, Stream } from 'effect';
import { constant } from 'effect/Function';
const tokens = Stream.fromIterable('Hello world').pipe(
Stream.schedule(Schedule.spaced('50 millis')),
Stream.map((char) => ({ _tag: 'token' as const, content: char })),
);
const heartbeat = Stream.tick('100 millis').pipe(
Stream.map(constant({ _tag: 'ping' as const })),
);
const merged = Stream.merge(tokens, heartbeat);
В downstream приходит:
{ _tag: 'token', content: 'H' }
{ _tag: 'ping' }
{ _tag: 'token', content: 'e' }
{ _tag: 'token', content: 'l' }
{ _tag: 'ping' }
...
Точный порядок зависит от тайминга. На реальной сети джиттер каждый раз другой, для детерминированного теста есть TestClock (13 · Testing).
Halt-стратегии
В коде выше есть ловушка. Stream.tick бесконечен. По умолчанию Stream.merge ждёт обоих, и весь merged тоже становится бесконечным. Когда токены закончатся, heartbeat будет тикать вечно, всё повиснет.
Лечится опцией haltStrategy:
'both'(по умолчанию), ждать, пока оба закончатся;'left', остановиться, как только закончится первый аргумент;'right', остановиться, как только закончится второй;'either', остановиться на любом из.
Для чат-бота нужен 'left': токены это сюжет, heartbeat это вспомогалка. Когда LLM договорил, heartbeat не нужен.
const merged = Stream.merge(tokens, heartbeat, { haltStrategy: 'left' });
Теперь как только последний токен выпущен, tokens закрывается, merge это видит, прерывает файбер heartbeat-а, и стрим чисто завершается. Cleanup-у heartbeat-а (clearInterval у callback) тоже прилетит interrupt, и он отработает.
Размеченный союз своими руками
Часто хочется не просто слить два стрима, а понимать на выходе, откуда пришёл элемент. В Effect был отдельный комбинатор mergeWithTag, который принимал объект { llm, system } и сам навешивал _tag. В v4 его убрали: он экономил ровно одну строку и при этом прятал типы. Пишем разметку явно, Stream.map перед Stream.merge:
import { Schedule, Stream } from 'effect';
const llm = Stream.make('token 1', 'token 2').pipe(
Stream.schedule(Schedule.spaced('100 millis')),
Stream.map((value) => ({ _tag: 'llm' as const, value })),
);
const system = Stream.make('connected', 'processing').pipe(
Stream.schedule(Schedule.spaced('150 millis')),
Stream.map((value) => ({ _tag: 'system' as const, value })),
);
const tagged = Stream.merge(llm, system);
// tagged: Stream<{ _tag: 'llm'; value: string } | { _tag: 'system'; value: string }>
В downstream приходят уже размеченные объекты:
{ _tag: 'llm', value: 'token 1' }
{ _tag: 'system', value: 'connected' }
{ _tag: 'llm', value: 'token 2' }
{ _tag: 'system', value: 'processing' }
Stream.merge в сигнатуре честно склеивает A | A2, поэтому union по тегам выводится сам, без единого ручного аннотирования. Разбирать его удобно через ts-pattern или Match. В Pulse это полезный паттерн: один стрим событий, в котором перемешаны успехи, провалы, ретраи, метрики.
Динамический список · mergeAll
Когда у тебя массив однотипных стримов и ты не знаешь их число статически, есть Stream.mergeAll:
import { Stream } from 'effect';
const streams: ReadonlyArray<Stream.Stream<number>> = [
Stream.make(1, 2),
Stream.make(10, 20),
Stream.make(100, 200),
];
const merged = Stream.mergeAll(streams, { concurrency: 'unbounded' });
Опции:
concurrency: 'unbounded', запустить все стримы одновременно;concurrency: N, держать в полёте максимумN, остальные ждут своей очереди;bufferSize: 16(по умолчанию), размер общей очереди выходных значений; producer-файберы backpressure-ются, когда очередь полна.
Ограничение: все стримы должны быть одного типа Stream<A, E, R>. TypeScript не выведет union по разным A. Если типы разные, приводи их к общему заранее (тем же Stream.map с тегом) или бери попарный Stream.merge, у которого union в сигнатуре.
Что взять с собой
Stream.merge(s1, s2, { haltStrategy }), важная опция в задачах “поток данных плюс бесконечный heartbeat”.- Размеченный союз собирается вручную:
Stream.mapс_tagна каждом источнике, потомStream.merge. КомбинатораmergeWithTagв v4 нет. mergeAll(streams, { concurrency, bufferSize }), динамический список, все стримы одного типа.bufferSizeэто общая выходная очередь; на её заполнении producer-файберы backpressure-ются.
Раздел 7 · broadcastN, broadcast, share
Зачем раздавать стрим
LLM-бот, ответ нужен сразу:
- клиенту по WebSocket;
- логгеру для дебага;
- в БД для истории чатов.
Запустить стрим три раза нельзя: это три вызова LLM, лишние деньги и лишнее время. Нужно запустить один раз и фан-аут в трёх consumer-ов. Это broadcasting, и под капотом он использует PubSub (см. 07 · Координация).
broadcastN, фиксированный набор подписчиков
Stream.broadcastN({ n, capacity }) создаёт N подписок до старта источника и возвращает кортеж из N стримов. Все они видят все элементы:
import { Array, Console, Effect, Random, Stream } from 'effect';
const source = Stream.fromIterableEffect(
Effect.all(Array.makeBy(5, () => Random.nextInt)),
);
const program = Effect.gen(function* () {
const [s1, s2] = yield* source.pipe(Stream.broadcastN({ n: 2, capacity: 8 }));
yield* Effect.zip(
s1.pipe(Stream.runForEach((n) => Console.log(`s1: ${n}`))),
s2.pipe(Stream.runForEach((n) => Console.log(`s2: ${n}`))),
{ concurrent: true },
);
});
await Effect.runPromise(Effect.scoped(program));
// s1: 1234, s2: 1234, s1: 5678, s2: 5678, ... ← одни и те же числа
Источник дёрнут один раз, оба consumer-а получили идентичные значения. broadcastN идеален, когда точно знаешь число подписчиков и они стартуют вместе.
broadcast, подписки на лету
Если consumer-ы появляются динамически (HTTP-клиенты в чате подключаются и отваливаются), нужна другая модель. Раньше она называлась broadcastDynamic; в v4 динамический вариант это и есть просто broadcast, а фиксированный уехал в broadcastN:
import { Console, Effect, Stream } from 'effect';
const program = Effect.gen(function* () {
const dynamic = yield* Stream.broadcast(source, { capacity: 'unbounded' });
// в одном fiber-е
yield* dynamic.pipe(Stream.runForEach((n) => Console.log(`fast: ${n}`)));
// в другом fiber-е, может быть позже
yield* dynamic.pipe(Stream.runForEach((n) => Console.log(`slow: ${n}`)));
});
broadcast стартует источник сразу при создании и держит его в работе, пока не закроется скоуп. Каждый запуск возвращённого стрима это новая подписка на тот же PubSub.
Подвох: поздние подписчики пропускают ранее опубликованные значения. Если consumer 1 уже получил 10 элементов, а consumer 2 подключился только сейчас, ему достанутся только 11-й и дальше.
replay, история для опоздавших
Чтобы поздние подписчики видели последние N элементов:
const dynamic = yield* Stream.broadcast(source, {
capacity: 10,
replay: 3,
});
replay: 3 это “храни последние 3 элемента для новых подписчиков”. Это скользящее окно размера N, не полная история.
Важное замечание: replay: 3 это минимум, не максимум. Гарантия “поздний подписчик получит хотя бы последние 3 элемента”. Если capacity больше и значения ещё не были вытеснены, ему может достаться больше. Чтобы строго ограничить, используй bounded capacity со sliding-стратегией.
share, lazy плюс reference counting
broadcast стартует источник сразу, даже если никто не подписан. И держит его в работе до закрытия скоупа, даже если все подписчики ушли. Для дорогих источников (открытый WebSocket, открытый файл, активный курсор БД) это плохо.
Stream.share решает обе проблемы:
import { Effect, Schedule, Stream } from 'effect';
const expensive = Stream.range(0, 1000).pipe(
Stream.tap(() => Effect.logInfo('expensive step')),
Stream.schedule(Schedule.spaced('10 millis')),
);
const program = Effect.gen(function* () {
const shared = yield* Stream.share(expensive, {
capacity: 10,
replay: 3,
idleTimeToLive: '5 seconds',
});
// ...здесь подключаются consumer-ы...
});
Что отличает share от broadcast:
- lazy start: источник стартует только когда подключится первый consumer;
- reference counting: пока есть хотя бы один активный consumer, источник работает;
- idle cleanup: когда последний consumer отвалился, запускается таймер
idleTimeToLive. Если за это время никто не подключился, источник останавливается, ресурсы освобождаются. Если кто-то подключился в течение таймера, источник продолжает крутиться.
Это идеальный примитив для “стрим стоит дорого, открывать заново каждый раз тяжело, но если никому не нужен полминуты, можно и закрыть”.
broadcastN vs broadcast vs share
| Параметр | broadcastN | broadcast | share |
|---|---|---|---|
| Кол-во подписчиков | фиксированное, N upfront | динамическое | динамическое |
| Старт источника | при создании | при создании | при первом подписчике |
| Поздние подписчики | их нет, все стартуют вместе | пропускают историю | replay плюс idle-warm |
| Когда источник стоп | при закрытии скоупа | при закрытии скоупа | через idleTimeToLive после последнего |
| Когда использовать | известное число consumer-ов | подписчики на лету | дорогой источник, ленивый старт |
Все три используют PubSub под капотом, и все три требуют Effect.scoped (или находиться в scoped-контексте), чтобы рантайм мог корректно прибрать PubSub и subscriber-файберы.
Подвохи
- Сайд-эффекты источника начинают происходить сразу (
broadcastN/broadcast) или с первой подписки (share), независимо от того, тянете ли вы значения downstream-ом. Если не хотите трогать источник до явного триггера, заворачивайте вshareс явным ленивым consumer-ом. replayимеет смысл только дляbroadcastиshare. УbroadcastNвсе подписки создаются до старта, “опоздавших” не бывает.- Без
Effect.scopedbroadcast*упадёт в типах: их сигнатура требуетScopeв контексте.
Что взять с собой
broadcastN({ n, capacity }), фиксированный набор consumer-ов, все стартуют вместе.broadcast(source, { capacity }), динамические подписчики, опоздавшие пропускают историю (или используйreplay).share(source, { idleTimeToLive }), ленивый старт и idle-cleanup, для дорогих источников.- Все три требуют скоупа и используют
PubSubвнутри.
Раздел 8 · Pulse · monitorEventsStream и пакетная запись в JSONL
Что у нас уже есть
MonitorEventsPubSubнаPubSub(07 · Координация) с типом события вида{ _tag: 'ProbeSuccess' | 'ProbeFailure' | 'SlaSwitched'; ... }.SlaнаTxRef(08 · Транзакции) с правилом “три провала переключают URL”.- Воркер, который тикает по целям мониторинга и публикует события в
MonitorEvents.
Нужно:
- Стрим событий, который можно дать UI, логгеру, и JSONL-писателю одновременно.
- JSONL-писатель: батчит события в файл, не пишет по одному (дисковая нагрузка), не теряет на падении.
- Опционально: UI с подпиской на последние N событий через
Stream.shareсreplay.
Шаг 1 · monitorEventsStream через Stream.fromPubSub
// pulse-<nick>/src/stream/monitor-stream.ts
import { Effect, Stream } from 'effect';
import { MonitorEventsPubSub } from '../concurrency/coordination.ts';
import type { MonitorEvent } from '../events.ts';
export const monitorEventsStream: Stream.Stream<MonitorEvent, never, MonitorEventsPubSub> =
Stream.unwrap(
Effect.gen(function* () {
const pubsub = yield* MonitorEventsPubSub;
return Stream.fromPubSub(pubsub);
}),
);
Тут два примитива работают в паре. Stream.fromPubSub(pubsub) подписывается на PubSub через скоуп самого стрима: каждый, кто вытягивает этот стрим, открывает свою подписку. Старые события (опубликованные до подписки) не приходят, как и положено PubSub.
Stream.unwrap нужен потому, что достать сам pubsub из контекста мы можем только внутри Effect. Он превращает “эффект, который отдаёт стрим” в стрим и сам вычитает Scope из требований. Раньше для scoped-случая существовал отдельный Stream.unwrapScoped; теперь их слили в один unwrap.
Заметь, что сервис тут вообще не нужен: monitorEventsStream это просто значение с MonitorEventsPubSub в требованиях. Стримы ленивы, оборачивать их в отдельный сервис есть смысл только когда у сервиса появляется собственное состояние.
Шаг 2 · пакетная запись через groupedWithin
// pulse-<nick>/src/stream/jsonl-writer.ts
import { Context, Effect, FileSystem, Layer, Stream } from 'effect';
import type { PlatformError } from 'effect/PlatformError';
import { MonitorEventsPubSub } from '../concurrency/coordination.ts';
import type { MonitorEvent } from '../events.ts';
import { monitorEventsStream } from './monitor-stream.ts';
export interface JsonlWriterImpl {
readonly run: Effect.Effect<void, PlatformError, MonitorEventsPubSub>;
}
export class JsonlWriter extends Context.Service<JsonlWriter, JsonlWriterImpl>()(
'our-pulse/JsonlWriter',
) {}
export const JsonlWriterLive = Layer.effect(
JsonlWriter,
Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const writeBatch = (batch: ReadonlyArray<MonitorEvent>) =>
fs.writeFile(
'events.jsonl',
new TextEncoder().encode(batch.map((e) => JSON.stringify(e)).join('\n') + '\n'),
{ flag: 'a' },
);
const run = monitorEventsStream.pipe(
Stream.groupedWithin(64, '1 second'),
Stream.tap(writeBatch),
Stream.runDrain,
);
return JsonlWriter.of({ run });
}),
);
Пара деталей по форме кода. Сервис объявляется через Context.Service<Self, Shape>()('id'), а Layer к нему пишется отдельно: это стандартная конвенция Effect 4 (04 · Services и Layer). FileSystem живёт прямо в ядре effect, отдельного пакета для платформенных сервисов больше нет.
Что мы получили:
groupedWithin(64, '1 second'), событие уходит на диск как только накопилось 64 штуки или прошла 1 секунда. На пике мы пишем большими блоками (дёшево для диска), в idle пишем хотя бы раз в секунду (важные события не теряются в RAM).Stream.tap(writeBatch), побочный эффект “записать в файл”, без преобразования стрима.Stream.runDrain, исполнить ради эффекта, выкинуть значения.
В MainLive запускаем run через Effect.forkScoped, и он живёт ровно столько, сколько живёт скоуп приложения. На runtime.dispose() finalizer-ы стрима закроют файл и подписку.
Шаг 3 · UI-подписка через Stream.share
В дашборде хочется показывать “последние 5 событий” сразу при открытии, плюс live-обновления. Подключается несколько вкладок? Все они должны делить одну подписку.
// pulse-<nick>/src/stream/events-feed.ts
import { Context, Effect, Layer, Stream } from 'effect';
import type { MonitorEvent } from '../events.ts';
import { monitorEventsStream } from './monitor-stream.ts';
export class EventsFeed extends Context.Service<EventsFeed, Stream.Stream<MonitorEvent>>()(
'our-pulse/EventsFeed',
) {}
export const EventsFeedLive = Layer.effect(
EventsFeed,
Stream.share(monitorEventsStream, {
capacity: 64,
replay: 5,
idleTimeToLive: '30 seconds',
}),
);
Что меняется:
- первый подписчик триггерит
share, который под капотом стартуетbroadcast, а тот подписывается наPubSubчерезmonitorEventsStream; - любая следующая вкладка подключается к той же подписке, не создавая лишних потребителей
PubSub-а; - поздняя вкладка видит последние 5 событий (
replay: 5), и дальше идёт обычная live-лента; - когда все вкладки закрылись, начинается отсчёт 30 секунд. Если за это время никто не вернулся, подписка на
PubSubзакрывается. Если кто-то открыл новую вкладку через 10 секунд, она цепляется к той же подписке.
Тут интересен сам Layer. Stream.share требует Scope, и раньше под такое был отдельный конструктор Layer.scoped. В v4 его нет: обычный Layer.effect сам вычитает Scope из требований и привязывает его к времени жизни Layer-а. Один конструктор на оба случая, и share живёт ровно столько, сколько живёт приложение (04 · Services и Layer).
Шаг 4 · агрегация в окне для метрик
Хочется метрику “успех/провал за последние 5 секунд”, печатать раз в 5 секунд:
import { Console, Stream } from 'effect';
import { monitorEventsStream } from './monitor-stream.ts';
const metricsLoop = monitorEventsStream.pipe(
Stream.groupedWithin(Number.MAX_SAFE_INTEGER, '5 seconds'),
Stream.tap((batch) => {
const successes = batch.filter((e) => e._tag === 'ProbeSuccess').length;
const failures = batch.filter((e) => e._tag === 'ProbeFailure').length;
return Console.log(`[5s] ok=${successes} fail=${failures}`);
}),
Stream.runDrain,
);
Трюк с Number.MAX_SAFE_INTEGER: мы не хотим срабатывание по числу, только по времени. groupedWithin отдаст всё, что накопилось за 5 секунд, и начнёт собирать следующее окно. Если за окно вообще ничего не пришло, отдаст пустой массив (это нормально, проверяй batch.length === 0, если хочешь не печатать).
И заметь, насколько чище стал код от того, что батч это обычный массив: filter и length работают прямо на нём, никакой конвертации из Chunk в массив и обратно.
Шаг 5 · тест на батчинг через TestClock
@effect/vitest плюс TestClock позволяют детерминированно проверить, что batch уходит ровно по таймеру:
import { describe, it } from '@effect/vitest';
import { Effect, Fiber, Ref, Stream } from 'effect';
import { TestClock } from 'effect/testing';
import { expect } from 'vitest';
describe('JsonlWriter batching', () => {
it.effect('flushes after 1 second even if buffer not full', () =>
Effect.gen(function* () {
const collected = yield* Ref.make<ReadonlyArray<number>>([]);
const source = Stream.make(1, 2, 3); // меньше 64
const program = source.pipe(
Stream.groupedWithin(64, '1 second'),
Stream.tap((batch) => Ref.update(collected, (xs) => [...xs, ...batch])),
Stream.runDrain,
);
const fiber = yield* Effect.forkChild(program);
// источник эмитнул всё мгновенно, но batch ещё не отдан
yield* TestClock.adjust('999 millis');
expect(yield* Ref.get(collected)).toEqual([]); // batch ждёт окна
yield* TestClock.adjust('1 millis');
yield* Fiber.await(fiber); // окно закрылось, batch вылетел
expect(yield* Ref.get(collected)).toEqual([1, 2, 3]);
}).pipe(Effect.provide(TestClock.layer())),
);
});
Три детали, на которых легко споткнуться. TestClock теперь живёт в подмодуле effect/testing и подключается как обычный Layer через TestClock.layer(); отдельного TestContext больше нет. Effect.fork переименован в Effect.forkChild, чтобы в самом имени было видно, что файбер привязан к родителю. И Fiber больше не yield-ится напрямую: дожидаемся через Fiber.await(fiber).
TestClock.adjust(d) сдвигает виртуальное время. Никакого await sleep(1000) в тесте, всё мгновенно и детерминировано. Это та же дисциплина, что в 13 · Testing.
Что взять с собой
Stream.unwrapплюсStream.fromPubSub(pubsub), базовый адаптерMonitorEventsPubSubвStream.Stream.groupedWithin(N, duration), batching для JSONL/БД/метрик. Окно по числу и по времени.Stream.shareплюсreplayиidleTimeToLive, идеальная подписка для UI с фан-аутом.TestClock, детерминированный тест на окна, никакихsleep-ов в тестах.
Финал · чек-лист
ДЗ
Дальше
Следующий урок · 10 · Batching и Request. Стримы хорошо ложатся на пакеты запросов: один шаг конвейера это пачка Effect.request, а RequestResolver сам решает, как их собрать и отправить. По сути это та же история, что groupedWithin плюс mapEffect, но завёрнутая в первоклассный примитив с дедупликацией и кешем.
Контекст по async-стримам в JS живёт в 05 · Async/06 Стримы и async iterators, там разобраны ReadableStream, Node-овский stream.pipeline, и WHATWG-протокол backpressure. Effect Stream это та же идея, поднятая в тип-систему. По координации см. 07 · Координация: Queue и PubSub это те самые примитивы, которые Effect Stream использует под капотом для callback и broadcast. По атомарности шага агрегации см. 08 · Транзакции: если внутри Stream.tap нужно поправить несколько счётчиков согласованно, заворачивай в Effect.tx.