../

Streaming data

Processing data piece by piece instead of all at once: WHATWG web streams (browser, Node, Bun, Deno, edge runtimes), Node's own streams, streaming HTTP bodies, Server-Sent Events and NDJSON, with the types for each. Files on disk are in File I/O.

Why stream

Buffered (await res.text())Streamed
Memorywhole payloadone chunk plus queues
Time to first byte usedafter the last byte arrivesas soon as the first chunk arrives
Slow consumerproducer fills memorybackpressure slows the producer
Cancel midwaydownload finishes anywaycancel() / abort stops the source
Codeone awaita pipeline of stages
  • Backpressure is end to end: a full queue in your code stops reads, the socket's receive buffer fills, and TCP flow control (the receive window) makes the sender pause. HTTP/2 adds its own per-stream windows.
  • HTTP/1.1 sends a body of unknown length with Transfer-Encoding: chunked; HTTP/2 and HTTP/3 frame it natively. Either way the server can start sending before it knows the total size.
  • Anything in the path that buffers (reverse proxies, CDNs, compression middleware) cancels the benefit; see Streaming HTTP responses.
  • Stream when data is large, unbounded (logs, events, LLM tokens) or slow to produce. For small JSON, buffer.

Web Streams

TypeRoleCreate with
ReadableStream<R>source of chunks of type Rnew ReadableStream(source), res.body, blob.stream()
WritableStream<W>sink that accepts Wnew WritableStream(sink)
TransformStream<I, O>writable side I, readable side Onew TransformStream(transformer)
ReadableStreamDefaultReader<R>locks a stream and reads itstream.getReader()
WritableStreamDefaultWriter<W>locks a sink and writessink.getWriter()
ReadableStreamBYOBReaderreads into your own buffergetReader({ mode: "byob" }) on a type: "bytes" stream

ReadableStream: start, pull, cancel

// pull source: called whenever the queue wants more
function counter(limit: number): ReadableStream<number> {
  let n = 0;
  return new ReadableStream<number>({
    pull(controller) {
      if (n >= limit) {
        controller.close();
        return;
      }
      controller.enqueue(n);
      n += 1;
    },
    cancel(reason) {
      // consumer gave up: release resources here
    },
  });
}
 
// push source: start() wires up events once
function ticks(ms: number): ReadableStream<number> {
  let id: ReturnType<typeof setInterval> | undefined;
  return new ReadableStream<number>({
    start(controller) {
      const tick = () => controller.enqueue(Date.now());
      id = setInterval(tick, ms);
    },
    cancel() {
      clearInterval(id);
    },
  });
}
Underlying source memberWhen it runs
start(controller)once, immediately; may return a promise
pull(controller)when desiredSize > 0; not called again until its promise settles
cancel(reason)consumer called cancel() or broke out of for await
type: "bytes"makes a byte stream (zero-copy BYOB reads)
Controller methodDoes
enqueue(chunk)adds a chunk to the queue
close()ends the stream after queued chunks are read
error(e)fails the stream; pending and future reads reject
desiredSizehigh-water mark minus queued size; <= 0 means stop producing

Reading with a reader

async function readAll(stream: ReadableStream<string>) {
  const reader = stream.getReader();
  const parts: string[] = [];
  try {
    while (true) {
      const { done, value } = await reader.read();
      if (done) break;
      parts.push(value);
    }
  } finally {
    reader.releaseLock();
  }
  return parts.join("");
}

WritableStream and TransformStream

const received: string[] = [];
const sink = new WritableStream<string>({
  write(chunk) {
    received.push(chunk); // may return a promise
  },
  close() {},
  abort(reason) {},
});
 
const writer = sink.getWriter();
await writer.ready; // wait for backpressure to clear
await writer.write("a");
await writer.close();
 
const upper = new TransformStream<string, string>({
  transform(chunk, controller) {
    controller.enqueue(chunk.toUpperCase());
  },
  flush(controller) {
    controller.enqueue("\n"); // after the last chunk
  },
});

Piping & backpressure

declare const source: ReadableStream<string>;
declare const upper: TransformStream<string, string>;
declare const sink: WritableStream<string>;
 
const ac = new AbortController();
await source
  .pipeThrough(upper) // returns upper.readable
  .pipeTo(sink, { signal: ac.signal }); // resolves when done
 
const [forLog, forUser] = source.tee(); // two readers
MethodReturnsNotes
pipeThrough(ts, opts?)ReadableStream<O>any { writable, readable } pair works
pipeTo(ws, opts?)Promise<void>rejects if either side errors
tee()[ReadableStream, ReadableStream]the slower branch's chunks queue without limit
cancel(reason)Promise<void>on an unlocked stream
lockedbooleantrue while a reader or pipe holds it
Pipe optionDefaultEffect
preventClosefalsekeep the sink open when the source ends
preventAbortfalsedo not abort the sink when the source errors
preventCancelfalsedo not cancel the source when the sink errors
signalnoneabort the pipe; errors both ends unless prevented

Queuing strategies

// count chunks (default: highWaterMark 1)
const byCount = new CountQueuingStrategy({
  highWaterMark: 16,
});
 
// count bytes: buffer up to 64 KiB before pausing the source
const byBytes = new ByteLengthQueuingStrategy({
  highWaterMark: 64 * 1024,
});
 
declare const source: UnderlyingDefaultSource<Uint8Array>;
const stream = new ReadableStream(source, byBytes);
 
// new TransformStream(transformer, writable, readable)
const buffered = new TransformStream({}, byCount, byCount);

desiredSize = highWaterMark - queued. A pull source gets backpressure for free; a push source (start plus events) must check controller.desiredSize and pause its producer itself. On the writing side, await writer.ready resolves when desiredSize is positive again.

Text & compression streams

async function fetchGzippedText(url: string) {
  const res = await fetch(url);
  if (!res.ok || !res.body) throw new Error(`${res.status}`);
  const text = res.body
    .pipeThrough(new DecompressionStream("gzip"))
    .pipeThrough(new TextDecoderStream()); // string chunks
  // not new Response(text): Response bodies must be bytes
  return (await Array.fromAsync(text)).join("");
}
 
async function gzip(text: string): Promise<Blob> {
  const compressed = new Blob([text])
    .stream()
    .pipeThrough(new CompressionStream("gzip"));
  return new Response(compressed).blob();
}
StreamInput to outputNotes
TextDecoderStream(label?)bytes to stringhandles characters split across chunks; fatal option
TextEncoderStream()string to Uint8ArrayUTF-8 only
CompressionStream(format)bytes to compressed bytesBaseline 2023
DecompressionStream(format)compressed bytes to byteserrors on corrupt input
FormatSupport
"gzip", "deflate"all engines, Node 18+, Bun
"deflate-raw"Baseline 2023
"brotli", "zstd"not cross-browser yet; in Node use node:zlib

Async iteration

declare const stream: ReadableStream<string>;
 
for await (const chunk of stream) {
  if (chunk === "STOP") break; // cancels the stream
}
 
// keep the stream alive after breaking out
declare const other: ReadableStream<string>;
for await (
  const c of other.values({ preventCancel: true })
) {
  break;
}
 
// collect everything
const all = await Array.fromAsync(other);
FeatureSupport
for await over ReadableStreamNode 16.5+, Deno, Bun, Chrome 124, Firefox 110, Safari 27
ReadableStream.from(iterable)Node 20.6+, Bun, Firefox 117, Safari 27; not in Chromium
TypeScriptiteration needs "lib": ["dom.asynciterable"]; ReadableStream.from is not in the DOM lib

Until ReadableStream.from lands everywhere, adapt an async generator by hand:

function fromIterable<T>(
  iterable: AsyncIterable<T> | Iterable<T>,
): ReadableStream<T> {
  const it =
    Symbol.asyncIterator in iterable
      ? iterable[Symbol.asyncIterator]()
      : iterable[Symbol.iterator]();
  return new ReadableStream<T>({
    async pull(controller) {
      const { value, done } = await it.next();
      if (done) controller.close();
      else controller.enqueue(value);
    },
    async cancel(reason) {
      await it.return?.(reason);
    },
  });
}
 
async function* pages(url: string) {
  for (let page = 1; page <= 3; page++) {
    const res = await fetch(`${url}?page=${page}`);
    yield (await res.json()) as unknown;
  }
}
 
const results = fromIterable(pages("/api/items"));

In Node, import { ReadableStream } from "node:stream/web" is typed with .from.

Node streams

ClassIsExamples
Readablesourcefs.createReadStream, process.stdin, HTTP request
Writablesinkfs.createWriteStream, process.stdout, HTTP response
Duplexindependent read and write sidesnet.Socket
Transformduplex whose output derives from inputzlib.createGzip(), crypto.createHash stream
PassThroughtransform that changes nothingtapping or joining pipelines
import {
  createReadStream,
  createWriteStream,
} from "node:fs";
import { Readable, Transform } from "node:stream";
import { pipeline } from "node:stream/promises";
 
const upper = new Transform({
  transform(chunk: Buffer, _encoding, callback) {
    callback(null, chunk.toString("utf8").toUpperCase());
  },
});
 
await pipeline(
  createReadStream("in.txt"),
  upper,
  createWriteStream("out.txt"),
);
 
// async generators work as pipeline stages
await pipeline(
  createReadStream("in.txt", { encoding: "utf8" }),
  async function* (source: AsyncIterable<string>) {
    for await (const text of source) yield text.trim();
  },
  createWriteStream("trimmed.txt"),
);
 
// object mode: chunks are any JS value
const rows = Readable.from([{ id: 1 }, { id: 2 }]);
for await (const row of rows) {
  // row: any; annotate or validate
}
DetailNode streams
Buffer thresholdhighWaterMark: 64 KiB of bytes (16 KiB before Node 22 and on Windows), 16 objects in object mode
Backpressurewrite() returns false; wait for "drain"
Object mode{ objectMode: true }; Readable.from() defaults to it
Errorsuse pipeline, not .pipe(): .pipe() neither forwards errors nor destroys on failure
Cleanupstream.destroy(err), finished(stream) from node:stream/promises

Converting between Node and web streams

import { createReadStream } from "node:fs";
import { Readable } from "node:stream";
import type {
  ReadableStream as NodeRS,
} from "node:stream/web";
 
declare const res: Response;
 
// web to Node: any async iterable works with Readable.from
const nodeA = Readable.from(res.body ?? []);
// or fromWeb; the cast bridges DOM vs node:stream/web types
const nodeB = Readable.fromWeb(
  res.body as unknown as NodeRS<Uint8Array>,
);
 
// Node to web, e.g. to return from a route handler
const web = Readable.toWeb(createReadStream("big.bin"));
const response = new Response(
  web as unknown as ReadableStream,
);
UseWhen
Web streamscode shared with browsers, edge runtimes, fetch, route handlers, Bun
Node streamsfs, zlib, net, child_process, and libraries that expect them
Bothconvert at the boundary with Readable.fromWeb / toWeb; keep one style inside a module

Streaming HTTP responses

Reading a body with progress

async function download(
  url: string,
  onProgress: (received: number, total?: number) => void,
): Promise<Blob> {
  const res = await fetch(url);
  if (!res.ok || !res.body) throw new Error(`${res.status}`);
  const length = Number(res.headers.get("content-length"));
  const total = length > 0 ? length : undefined;
 
  const chunks: Uint8Array<ArrayBuffer>[] = [];
  let received = 0;
  for await (const chunk of res.body) {
    chunks.push(chunk);
    received += chunk.byteLength;
    onProgress(received, total);
  }
  return new Blob(chunks);
}

Content-Length is the size on the wire: for a compressed response it is smaller than the decoded bytes you count, so clamp progress at 100%.

Returning a stream from a route handler

app/api/words/route.ts
const encoder = new TextEncoder();
const sleep = (ms: number) =>
  new Promise((resolve) => setTimeout(resolve, ms));
 
export function GET(req: Request): Response {
  const stream = new ReadableStream<Uint8Array>({
    async start(controller) {
      for (const word of ["streams", "are", "neat"]) {
        if (req.signal.aborted) break; // client left
        controller.enqueue(encoder.encode(`${word}\n`));
        await sleep(300);
      }
      controller.close();
    },
  });
  return new Response(stream, {
    headers: {
      "Content-Type": "text/plain; charset=utf-8",
      "X-Accel-Buffering": "no", // nginx: do not buffer
    },
  });
}

The same (req: Request) => Response shape works in Next.js route handlers, Bun, Deno, Hono and Cloudflare Workers.

Bun.serve

declare function makeStream(): ReadableStream<Uint8Array>;
 
const server = Bun.serve({
  port: 3000,
  idleTimeout: 0, // seconds; default 10 closes quiet streams
  routes: {
    "/stream": () => new Response(makeStream()),
  },
  fetch() {
    return new Response("Not found", { status: 404 });
  },
});
Buffering culpritFix
nginx and similar proxiesX-Accel-Buffering: no, or proxy_buffering off
Compression middlewareflush per chunk, or skip compression for streams
CDNscheck whether the plan passes chunked responses through
Safaribuffers the first 1024 bytes of a response before rendering
curluse curl -N; it still flushes only on newlines
Uploadsstreaming request bodies need duplex: "half" and are Chromium-only over HTTP/2

Server-Sent Events

One-way, server-to-client events over a plain HTTP response with Content-Type: text/event-stream. The browser parses it with EventSource, which reconnects automatically.

retry: 5000
: comment lines start with a colon (use as heartbeat)
 
id: 41
event: price
data: {"sym":"ACME","px":12.5}
 
id: 42
data: first line
data: second line
 
FieldMeaning
data:payload; several data: lines are joined with \n
event:event name for addEventListener; default "message"
id:stored as the last event id; sent back on reconnect
retry:reconnect delay in milliseconds
:comment, ignored; keeps idle connections open
blank linedispatches the event

Typed EventSource client

type ServerEvents = {
  price: { sym: string; px: number };
  notice: { text: string };
};
 
function on<K extends keyof ServerEvents>(
  source: EventSource,
  type: K,
  handler: (data: ServerEvents[K]) => void,
): () => void {
  const listener = (e: MessageEvent<string>) => {
    // validate with Zod in real code; the wire is untrusted
    handler(JSON.parse(e.data) as ServerEvents[K]);
  };
  source.addEventListener(type, listener);
  return () => source.removeEventListener(type, listener);
}
 
const source = new EventSource("/api/prices", {
  withCredentials: false, // true sends cookies cross-origin
});
const off = on(source, "price", ({ sym, px }) => {});
source.onerror = () => {
  // readyState CONNECTING: retrying; CLOSED: gave up
};
// later: off(); source.close();

Server

app/api/prices/route.ts
const encoder = new TextEncoder();
 
function sseMessage(
  data: unknown,
  opts: { event?: string; id?: string } = {},
): Uint8Array {
  const lines = [
    opts.event && `event: ${opts.event}`,
    opts.id && `id: ${opts.id}`,
    `data: ${JSON.stringify(data)}`,
  ].filter(Boolean);
  return encoder.encode(`${lines.join("\n")}\n\n`);
}
 
export function GET(req: Request): Response {
  const lastId = Number(
    req.headers.get("last-event-id") ?? 0,
  );
  let timer: ReturnType<typeof setInterval> | undefined;
 
  const stream = new ReadableStream<Uint8Array>({
    start(controller) {
      let id = lastId; // resume after the last one seen
      timer = setInterval(() => {
        id += 1;
        const msg = sseMessage(
          { sym: "ACME", px: 10 + Math.random() },
          { event: "price", id: String(id) },
        );
        controller.enqueue(msg);
      }, 1000);
    },
    cancel() {
      clearInterval(timer); // client disconnected
    },
  });
 
  return new Response(stream, {
    headers: {
      "Content-Type": "text/event-stream",
      "Cache-Control": "no-cache",
      "X-Accel-Buffering": "no",
    },
  });
}
BehaviorDetail
Reconnectautomatic after a network drop, waiting retry ms (browser default is a few seconds)
Last-Event-IDrequest header on reconnect carrying the last id:; replay from there
Stop reconnectingrespond with a non-200 status or another content type; client calls close()
HeadersEventSource cannot set custom headers; use cookies, or a fetch-based client for auth headers
Connection limitabout 6 per origin on HTTP/1.1 across all tabs; serve over HTTP/2
Binarytext only (UTF-8); base64 binary or use WebSockets

SSE or WebSockets

SSEWebSockets
Directionserver to clientboth ways
Transportordinary HTTP; works with proxies, HTTP/2, auth cookiesupgraded TCP connection
Reconnect and resumebuilt in (retry, Last-Event-ID)write it yourself
DataUTF-8 texttext and binary
Fitsnotifications, feeds, progress, LLM token streamschat, games, collaborative editing

NDJSON & line parsing

Newline-delimited JSON (application/x-ndjson, also called JSON Lines): one JSON value per line, so each line parses on its own as it arrives.

function splitLines(): TransformStream<string, string> {
  let buffer = "";
  return new TransformStream<string, string>({
    transform(chunk, controller) {
      buffer += chunk;
      const lines = buffer.split(/\r?\n/);
      buffer = lines.pop() ?? ""; // keep the partial line
      for (const line of lines) controller.enqueue(line);
    },
    flush(controller) {
      if (buffer) controller.enqueue(buffer);
    },
  });
}
import { z } from "zod";
 
declare function splitLines(): TransformStream<
  string,
  string
>;
 
function parseNdjson<S extends z.ZodType>(
  schema: S,
): TransformStream<string, z.infer<S>> {
  return new TransformStream<string, z.infer<S>>({
    transform(line, controller) {
      if (line.trim() === "") return; // skip blank lines
      const result = schema.safeParse(JSON.parse(line));
      if (result.success) controller.enqueue(result.data);
      else controller.error(result.error);
    },
  });
}
 
const LogEntry = z.object({
  level: z.enum(["info", "warn", "error"]),
  msg: z.string(),
});
 
const res = await fetch("/api/logs.ndjson");
if (!res.body) throw new Error("No body");
const entries = res.body
  .pipeThrough(new TextDecoderStream())
  .pipeThrough(splitLines())
  .pipeThrough(parseNdjson(LogEntry));
 
for await (const entry of entries) {
  if (entry.level === "error") console.error(entry.msg);
}

Writing NDJSON is JSON.stringify(value) + "\n" per record (JSON.stringify never emits a raw newline). In Node, readline over a file stream splits lines too; see File I/O.

Typing streams

TypeWhere it comes from
ReadableStream<Uint8Array<ArrayBuffer>>res.body (nullable), blob.stream()
ReadableStream<string>after pipeThrough(new TextDecoderStream())
ReadableStreamDefaultReader<T>stream.getReader()
ReadableStreamReadResult<T>reader.read(): { done: false, value: T } | { done: true, value: undefined }
UnderlyingDefaultSource<T>the object passed to new ReadableStream
Transformer<I, O>the object passed to new TransformStream
UnderlyingSink<W>the object passed to new WritableStream
QueuingStrategy<T>{ highWaterMark, size(chunk) }
ReadableWritablePair<O, I>anything pipeThrough accepts

Generic helpers

function map<I, O>(
  fn: (chunk: I) => O | Promise<O>,
): TransformStream<I, O> {
  return new TransformStream<I, O>({
    async transform(chunk, controller) {
      controller.enqueue(await fn(chunk));
    },
  });
}
 
function filter<T, S extends T>(
  pred: (chunk: T) => chunk is S,
): TransformStream<T, S> {
  return new TransformStream<T, S>({
    transform(chunk, controller) {
      if (pred(chunk)) controller.enqueue(chunk);
    },
  });
}
 
async function first<T>(
  stream: ReadableStream<T>,
): Promise<T | undefined> {
  const reader: ReadableStreamDefaultReader<T> =
    stream.getReader();
  try {
    const result = await reader.read();
    return result.done ? undefined : result.value;
  } finally {
    await reader.cancel(); // done with the rest
  }
}
 
declare const numbers: ReadableStream<number | null>;
const doubled = numbers
  .pipeThrough(filter((n): n is number => n !== null))
  .pipeThrough(map((n) => n * 2)); // ReadableStream<number>

Recipes

Render streamed text into the page

When showing an LLM answer or a long log as it arrives, without ever parsing it as HTML.

async function streamInto(
  el: HTMLElement,
  url: string,
  body: unknown,
  signal?: AbortSignal,
): Promise<void> {
  const res = await fetch(url, {
    method: "POST",
    headers: { "Content-Type": "application/json" },
    body: JSON.stringify(body),
    signal, // "Stop" button: controller.abort()
  });
  if (!res.ok || !res.body) throw new Error(`${res.status}`);
  const reader = res.body
    .pipeThrough(new TextDecoderStream())
    .getReader();
  el.textContent = "";
  for (;;) {
    const { done, value } = await reader.read();
    if (done) break;
    el.append(value); // a text node: never parsed as HTML
  }
}

Stream a large export as NDJSON

When a download is too big to build in memory and the database cursor should only advance as fast as the client reads.

app/api/export/route.ts
type Row = { id: number; email: string };
declare function queryRows(): AsyncIterable<Row>; // cursor
 
const encoder = new TextEncoder();
 
export function GET(): Response {
  const rows = queryRows()[Symbol.asyncIterator]();
  const body = new ReadableStream<Uint8Array>({
    async pull(controller) { // only when the client reads
      const { done, value } = await rows.next();
      if (done) controller.close();
      else controller.enqueue(
        encoder.encode(`${JSON.stringify(value)}\n`),
      );
    },
    async cancel() {
      await rows.return?.(); // client left: close the cursor
    },
  });
  return new Response(body, {
    headers: { "Content-Type": "application/x-ndjson" },
  });
}

SSE with auth headers over fetch

When an event stream needs an Authorization header, which EventSource can't send; splitLines() is from NDJSON & line parsing.

type SseEvent = { event: string; data: string };
 
async function* sse(
  url: string,
  init: RequestInit, // e.g. headers: { Authorization }
): AsyncGenerator<SseEvent> {
  const res = await fetch(url, init);
  if (!res.ok || !res.body) throw new Error(`${res.status}`);
  const lines = res.body
    .pipeThrough(new TextDecoderStream())
    .pipeThrough(splitLines());
  let event = "message";
  let buf: string[] = [];
  for await (const line of lines) {
    if (line === "") { // blank line ends an event
      if (buf.length) yield { event, data: buf.join("\n") };
      [event, buf] = ["message", []];
      continue;
    }
    const [field, ...rest] = line.split(":");
    const value = rest.join(":").replace(/^ /, "");
    if (field === "data") buf = [...buf, value];
    else if (field === "event") event = value;
  }
}

There is no automatic reconnect, id or retry handling here; wrap the loop in the backoff from WebSockets if you need it.

Batch chunks

When the sink is cheaper per batch than per item (bulk inserts, batch API endpoints).

function batch<T>(size: number): TransformStream<T, T[]> {
  let buffer: T[] = [];
  return new TransformStream<T, T[]>({
    transform(item, controller) {
      buffer.push(item);
      if (buffer.length >= size) {
        controller.enqueue(buffer);
        buffer = [];
      }
    },
    flush(controller) {
      if (buffer.length > 0) controller.enqueue(buffer);
    },
  });
}
 
type Row = { id: number };
declare const rows: ReadableStream<Row>;
declare function insertMany(rows: Row[]): Promise<void>;
 
// one INSERT per 500 rows; backpressure waits for the DB
await rows.pipeThrough(batch<Row>(500)).pipeTo(
  new WritableStream({ write: (rs) => insertMany(rs) }),
);

Hash a large file

When checksumming uploads or build artifacts without loading the file into memory.

import { createHash } from "node:crypto";
import { createReadStream } from "node:fs";
 
async function sha256File(path: string): Promise<string> {
  const hash = createHash("sha256");
  for await (const chunk of createReadStream(path)) {
    hash.update(chunk as Buffer); // constant memory
  }
  return hash.digest("hex");
}
 
// Bun: native hasher over a web stream
async function sha256Bun(path: string): Promise<string> {
  const hasher = new Bun.CryptoHasher("sha256");
  for await (const chunk of Bun.file(path).stream()) {
    hasher.update(chunk);
  }
  return hasher.digest("hex");
}

References