Площадка для изучения стримов в Node.js и браузере — от буферов и pipe() до параллельной обработки через worker threads и потокового аудио.
Стек: Node.js 24+, TypeScript 5.8, ESM
Требования: Node.js >= 24, yarn
git clone <repo-url>
cd streams-playground
yarn installКак запускать примеры:
- Standalone-файлы:
npx tsx <файл.ts>илиnode <файл.js> - Проекты с
package.json: перейти в директорию и выполнитьnpm start
| # | Подпроект | Тема |
|---|---|---|
| 01 | ecommerce-app | Observer-паттерн на EventEmitter |
| 02 | buffer | Основы Buffer: аллокация, fill, кодировки |
| 03 | buffer-streams | Ручная разбивка файла на чанки |
| 04 | readable-writable-transform | ETL-пайплайн, чтение больших файлов |
| 05 | duplex-streams | Duplex, Transform, PassThrough, broadcast |
| 06 | net-chat | TCP-чат на net.Socket |
| 07 | pipe-vs-pipeline | Разница pipe() и pipeline() при обрыве стрима |
| 08 | csv-to-ndjson | CSV → NDJSON конвертер с прогрессом |
| 09 | async-iterators | Async-генераторы вместо классов стримов |
| 10 | stream-operators-with-sqlite | Операторы .map()/.filter() + SQLite |
| 11 | abort-controller-with-stream | Отмена стрима через AbortController |
| 12 | abort-controller-with-server | Abort/resume HTTP-стрима |
| # | Подпроект | Тема |
|---|---|---|
| 13 | webstreams-introduction | ReadableStream, WritableStream, TransformStream |
| 14 | csv-streaming-browser | Потоковая загрузка CSV в браузер |
| # | Подпроект | Тема |
|---|---|---|
| 01 | child-process | Фильтрация NDJSON через child_process.fork() |
| 02 | worker-threads | То же, но через worker_threads |
| 03 | large-report-in-browser | Web Workers + TransformStream в браузере |
| 04 | spotify-radio | Потоковое аудио, микширование через SoX |
Реализация паттерна Observer/Pub-Sub. PaymentSubject хранит подписчиков в Set и рассылает уведомления при оплате. Подписчики — Marketing (отправляет welcome-email) и Shipment (упаковывает заказ). После unsubscribe уведомления перестают приходить.
Файлы: src/subjects/paymentSubject.js, src/events/payment.js, src/observers/marketing.js, src/observers/shipment.js
cd 1-streams-in-node-js/01-ecommerce-app
npm start
npm test # Jest, coverage 100%Работа с Buffer на низком уровне: Buffer.alloc(3), fill('hello') (обрезается до 'hel'), вставка через set(). Конвертация строки "Hello there" в массив char-кодов и hex-байтов с обратным преобразованием.
Файлы: index.js
cd 1-streams-in-node-js/02-buffer
node index.jsЗачем нужны стримы: файл из 20 строк читается целиком в память, затем вручную разбивается на чанки по 4 строки (5 итераций). Для каждого чанка создаётся Buffer и замеряется byteLength. Демонстрирует неудобство ручного чанкинга — мотивация перейти к настоящим стримам.
Файлы: index.js, file.txt
cd 1-streams-in-node-js/03-buffer-streams
node index.jsДва примера:
ETL.js — генерирует 1 000 000 записей (UUID + имя), прогоняет через два Transform-стрима (mapFields → CSV + uppercase, mapHeaders → добавляет заголовок) и пишет в my.csv (~48 МБ).
read-big-file.js — показывает, что fs.readFile() падает на файле > 2 ГБ (ERR_FS_FILE_TOO_LARGE), а createReadStream() читает его чанками по 65 КБ без проблем.
Файлы: ETL.js, read-big-file.js, my.csv
cd 1-streams-in-node-js/04-readable-writable-transform-streams
node ETL.js
node read-big-file.jsduplex-transform-api.js — Duplex-стрим с независимыми read/write каналами: read() генерирует данные раз в секунду, write() логирует входящие. Transform-стрим transformToUpperCase показывает, что Transform — это Duplex, где выход зависит от входа. Финал: кольцевой пайпинг duplexServer → transform → duplexServer → outputStream.
duplex-broadcast.ts — PassThrough как broadcaster: читает из файла и одновременно пишет в несколько клиентских стримов (output-{uuid}.txt).
Файлы: duplex-transform-api.js, duplex-broadcast.ts
cd 1-streams-in-node-js/05-duplex-streams
node duplex-transform-api.js
npx tsx duplex-broadcast.tsTCP-чат на net.createServer(). Сервер хранит сокеты в Map с UUID, при получении сообщения рассылает его остальным клиентам (broadcast). Клиент использует PassThrough + readline.cursorTo() для чистого терминального интерфейса.
Порт: 3000
Файлы: 01.server.ts, 02.client.ts, extend.d.ts
cd 1-streams-in-node-js/06-net-chat
# Терминал 1
npx tsx 01.server.ts
# Терминал 2, 3, ...
npx tsx 02.client.tsДва HTTP-сервера показывают разницу:
- pipe() (порт 3000) — при досрочном
destroy()стрима ошибки не будет - pipeline() (порт 3001) — выбросит
ERR_STREAM_PREMATURE_CLOSE
Скрипт запускает оба сервера, делает запросы и уничтожает стримы, наглядно показывая поведение.
Файлы: index.mts
cd 1-streams-in-node-js/07-pipe-vs-pipeline
npx tsx index.mtsКонвертер CSV → NDJSON с отображением прогресса. Пайплайн: createReadStream → CSVToNDJSON (кастомный Transform, буферизует неполные строки) → processData (добавляет инкрементальные ID) → reporter.progress() (PassThrough, считает процент) → createWriteStream.
Файлы: src/index.ts, src/stream-components/csvtondjson.ts, src/stream-components/reporter.ts, src/util.ts
cd 1-streams-in-node-js/08-csv-to-ndjson
npm start
npm test # Vitest
npm run test:coveragepipeline() принимает не только классы стримов, но и async-генераторы. Четыре функции — readable, transform, duplex, writable — реализованы как async function* и связаны через pipeline(). Transform заменяет пробелы на _, duplex агрегирует данные и считает байты.
Файлы: index.ts
cd 1-streams-in-node-js/09-async-iterators
npx tsx index.tsПотоковое чтение из SQLite с пагинацией (LIMIT 100 / OFFSET). Async-генератор selectAsStream() отдаёт строки по одной. Цепочка .map() операторов: первый — uppercase имени (async), второй — сериализация в NDJSON. Результат пишется в data/output.ndjson.
Схема БД: users (id TEXT, name TEXT, age NUMBER, company TEXT)
Файлы: src/index.ts, src/seed.ts
cd 1-streams-in-node-js/10-stream-operators-with-sqlite
npm run seed # Создать БД и заполнить 100 записями (faker.js)
npm startAbortController для отмены pipeline(). Readable-генератор отдаёт тики каждые 200 мс. Через 500 мс вызывается abortController.abort(), pipeline выбрасывает ABORT_ERR, который перехватывается в catch. Событие signal.onabort срабатывает при отмене.
Файлы: index.ts
cd 1-streams-in-node-js/11-abort-controller-with-stream
npx tsx index.tsHTTP-сервер стримит 600 JSON-объектов (один каждые 40 мс). Клиент использует Fetch API с AbortController.signal и демонстрирует паттерн abort/resume: отмена на 500 мс, переподключение на 1000 мс, повторная отмена на 1500 мс.
Порт: 3000
Файлы: src/server.ts, src/client.ts
cd 1-streams-in-node-js/12-abort-controller-with-server
# Терминал 1
npm run server
# Терминал 2
npm run clientWeb Streams API из node:stream/web: ReadableStream генерирует таймстамп-сообщения каждые 200 мс, TextDecoderStream декодирует буферы в строки, TransformStream приводит к верхнему регистру, WritableStream выводит в stdout. Всё связано через .pipeThrough() и .pipeTo().
Файлы: index.ts
cd 2-streams-in-browser/13-webstreams-introduction
npx tsx index.tsПолностековое приложение. Сервер (порт 3000) стримит CSV-файл с аниме (animeflv.csv), преобразуя в NDJSON через csvtojson. Браузерный клиент получает поток через Fetch API, парсит NDJSON по чанкам (обработка разрывов на границе чанка) и динамически рендерит карточки. Кнопки Start/Stop управляют потоком через AbortController.
Файлы: server/index.ts, client/index.html, client/index.js, client/index.test.js
# Терминал 1 — сервер
cd 2-streams-in-browser/14-csv-streaming-browser/server
npm run dev
# Терминал 2 — клиент
cd 2-streams-in-browser/14-csv-streaming-browser/client
npm start # http-server
npm test # Vitest (jsdom)child_process.fork() спавнит несколько процессов, каждый читает свой NDJSON-файл и фильтрует записи с Gmail-адресами. Результаты передаются в родительский процесс через IPC (process.send()). PassThrough-стрим объединяет (merge) читаемые стримы от всех дочерних процессов в один выходной файл output-gmail.ndjson.
Файлы: src/index.ts, src/background.job.ts
cd 3-streams-using-parallelism/01-child-process
npm startТа же задача (фильтрация Gmail из NDJSON), но через worker_threads — эффективнее за счёт разделяемой памяти. Три реализации воркера:
background.job.ts—.pipe()+ event listenersbackground.job-using-pipeline.ts—pipeline()+ async-генераторыbackground.job-using-readline.ts—readline.createInterface()+ async iteration
По умолчанию используется вариант с pipeline.
Файлы: index.ts, background.job.ts, background.job-using-pipeline.ts, background.job-using-readline.ts
cd 3-streams-using-parallelism/02-worker-threads
npm run build
npm startMVC-архитектура в браузере. Пользователь загружает CSV-файл и вводит поисковый запрос (regex). Обработка идёт через кастомный TransformStream (CSV → JSON). Чекбокс переключает между выполнением в основном потоке и Web Worker — наглядно видно, как Worker не блокирует UI. Прогресс-бар обновляется в реальном времени.
Файлы: index.html, src/index.js, src/controller.js, src/service.js, src/view.js, src/worker.js
cd 3-streams-using-parallelism/03-large-report-in-browser-worker-threads
npm start # browser-syncПотоковое интернет-радио. Сервер стримит MP3 (conversation.mp3) всем подключённым клиентам через PassThrough-стримы (каждый клиент получает свой, идентифицированный по UUID). Микширование звуковых эффектов (аплодисменты, бу, смех) происходит на лету через child_process.spawn() + SoX. Throttle-transform контролирует битрейт. Громкость: песня 0.99, эффекты 0.05.
Порт: 3000
Маршруты: /home — слушать, /controller — управлять (start/stop/fx), /stream — аудиопоток
Файлы: server/index.js, server/service.js, server/controller.js, server/routes.js, server/config.js
cd 3-streams-using-parallelism/04-spotify-radio
npm start
# Или через Docker
docker-compose upТребует установленного SoX в системе.
| Проект | Фреймворк | Команда |
|---|---|---|
| 01-ecommerce-app | Jest | npm test (coverage 100%) |
| 08-csv-to-ndjson | Vitest | npm test, npm run test:coverage |
| 14-csv-streaming-browser/client | Vitest (jsdom) | npm test |
| 04-spotify-radio | Node.js test runner | npm test, npm run test:watch |
- Readable / Writable / Transform / Duplex / PassThrough — все типы стримов с примерами
- pipe() vs pipeline() —
pipeline()корректно разрушает стримы и пробрасывает ошибки,pipe()— нет - Backpressure — автоматическое управление скоростью через
highWaterMark - Async-генераторы —
async function*как замена классам стримов вpipeline() - Web Streams API —
ReadableStream,WritableStream,TransformStream,pipeThrough(),pipeTo() - AbortController — отмена и возобновление потоковых операций через
signal - Stream operators —
.map(),.filter()как функциональный интерфейс стримов - Параллельная обработка —
child_process.fork(),worker_threads, Web Workers - Observer / Pub-Sub — EventEmitter + подписка/отписка наблюдателей