Skip to content

Skippia/streams-playground

Repository files navigation

streams-playground

Площадка для изучения стримов в 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

📂 Структура проекта

Часть 1 — Стримы в Node.js

# Подпроект Тема
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-стрима

Часть 2 — Стримы в браузере

# Подпроект Тема
13 webstreams-introduction ReadableStream, WritableStream, TransformStream
14 csv-streaming-browser Потоковая загрузка CSV в браузер

Часть 3 — Параллелизм

# Подпроект Тема
01 child-process Фильтрация NDJSON через child_process.fork()
02 worker-threads То же, но через worker_threads
03 large-report-in-browser Web Workers + TransformStream в браузере
04 spotify-radio Потоковое аудио, микширование через SoX

📖 Часть 1: Стримы в Node.js

01 — ecommerce-app (Observer)

Реализация паттерна 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%

02 — buffer

Работа с 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

03 — buffer-streams

Зачем нужны стримы: файл из 20 строк читается целиком в память, затем вручную разбивается на чанки по 4 строки (5 итераций). Для каждого чанка создаётся Buffer и замеряется byteLength. Демонстрирует неудобство ручного чанкинга — мотивация перейти к настоящим стримам.

Файлы: index.js, file.txt

cd 1-streams-in-node-js/03-buffer-streams
node index.js

04 — readable-writable-transform

Два примера:

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.js

05 — duplex-streams

duplex-transform-api.js — Duplex-стрим с независимыми read/write каналами: read() генерирует данные раз в секунду, write() логирует входящие. Transform-стрим transformToUpperCase показывает, что Transform — это Duplex, где выход зависит от входа. Финал: кольцевой пайпинг duplexServer → transform → duplexServer → outputStream.

duplex-broadcast.tsPassThrough как 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.ts

06 — net-chat

TCP-чат на 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

07 — pipe vs pipeline

Два 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

08 — csv-to-ndjson

Конвертер CSV → NDJSON с отображением прогресса. Пайплайн: createReadStreamCSVToNDJSON (кастомный 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:coverage

09 — async-iterators

pipeline() принимает не только классы стримов, но и 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

10 — stream-operators-with-sqlite

Потоковое чтение из 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 start

11 — abort-controller-with-stream

AbortController для отмены 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.ts

12 — abort-controller-with-server

HTTP-сервер стримит 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 client

🌐 Часть 2: Стримы в браузере

13 — webstreams-introduction

Web 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

14 — csv-streaming-browser

Полностековое приложение. Сервер (порт 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)

⚡ Часть 3: Параллелизм

01 — child-process

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

02 — worker-threads

Та же задача (фильтрация Gmail из NDJSON), но через worker_threads — эффективнее за счёт разделяемой памяти. Три реализации воркера:

  • background.job.ts.pipe() + event listeners
  • background.job-using-pipeline.tspipeline() + async-генераторы
  • background.job-using-readline.tsreadline.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 start

03 — large-report-in-browser

MVC-архитектура в браузере. Пользователь загружает 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

04 — spotify-radio

Потоковое интернет-радио. Сервер стримит 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 APIReadableStream, WritableStream, TransformStream, pipeThrough(), pipeTo()
  • AbortController — отмена и возобновление потоковых операций через signal
  • Stream operators.map(), .filter() как функциональный интерфейс стримов
  • Параллельная обработкаchild_process.fork(), worker_threads, Web Workers
  • Observer / Pub-Sub — EventEmitter + подписка/отписка наблюдателей

About

Playgroud to play around with streams in Node.js

Topics

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

No releases published

Packages

 
 
 

Contributors