Zum Inhalt springen

Große Datenmengen in Node.js streamen mit Backpressure

Schluss mit dem Laden ganzer Ergebnismengen: wie Node.js-Streams und asynchrone Generatoren Millionen Datensätze verarbeiten, ohne den Heap zu sprengen.

5 Min. Lesezeit
Diagramm einer Node.js-Stream-Pipeline, das zeigt, wie Daten von einer Datenbank über eine Transformationsstufe zu einer HTTP-Antwort fließen

Die meisten APIs, die „in der Entwicklung einwandfrei funktionieren", tragen eine stille Zeitbombe in sich: Sie laden komplette Ergebnismengen in den Speicher, bevor sie überhaupt etwas damit tun. Bei tausend Zeilen fällt das nicht auf. Bei hunderttausend schießen die Antwortzeiten in die Höhe. Bei einer Million läuft der Prozess in ein OOM und reißt den Pod mit sich. Streams sind die Standardlösung — aber gerade bei den Implementierungsdetails, vor allem beim Backpressure, scheitern die meisten Versuche.

Warum das Laden von allem die Produktion lahmlegt

Das problematische Muster begegnet einem überall:

tstypescript
// ❌ Loads the entire table into a JS array before the first byte is sent
export async function GET(req: Request) {
  const users = await db.query<User>(
    "SELECT * FROM users WHERE tenant_id = $1",
    [req.tenantId],
  );
  return Response.json(users); // serializes all rows at once — heap spikes with result size
}
 
// ✅ Streams rows as they arrive — heap usage stays flat regardless of row count
export async function GET(req: Request) {
  const rows = db.queryStream<User>(
    "SELECT * FROM users WHERE tenant_id = $1",
    [req.tenantId],
  );
  return new Response(serializeNDJSON(rows), {
    headers: { "Content-Type": "application/x-ndjson" },
  });
}

Die Streaming-Variante verarbeitet jeweils eine Zeile. Ob du 500 Zeilen sendest oder 5 Millionen, der Heap-Verbrauch bleibt nahezu konstant. Der Preis dafür ist Komplexität — und die steckt größtenteils im Backpressure.

Die Stream-Abstraktion in Node.js

Node.js-Streams gibt es in vier Varianten: Readable, Writable, Duplex und Transform. Für Pipelines mit großen Datenmengen arbeitest du meist mit Readable (der Datenquelle) und Transform (der Datenformung), die per Pipe in einen Writable (HTTP-Antwort, Datei, Message Queue) münden.

Der kanonische Verbindungspunkt ist stream.pipeline, nicht .pipe(). Der entscheidende Unterschied: pipeline propagiert Fehler und räumt bei einem Fehlschlag alle beteiligten Streams auf. .pipe() hinterlässt hängende Event-Listener, wenn mitten im Stream ein Fehler auftritt.

tstypescript
import { pipeline, Transform } from "node:stream";
import { promisify } from "node:util";
 
const pipelineAsync = promisify(pipeline);
 
async function exportUsersToResponse(
  queryStream: NodeJS.ReadableStream,
  res: NodeJS.WritableStream,
): Promise<void> {
  const ndJsonTransform = new Transform({
    objectMode: true,
    transform(row: unknown, _encoding, callback) {
      try {
        this.push(JSON.stringify(row) + "\n");
        callback();
      } catch (err) {
        callback(err as Error);
      }
    },
  });
 
  // All three streams are destroyed together on error or completion
  await pipelineAsync(queryStream, ndJsonTransform, res);
}

Jeder Stream in der Pipeline wird aufgeräumt, sobald irgendeine Stufe einen Fehler wirft. Keine manuellen .destroy()-Aufrufe, verstreut über deine Handler.

Backpressure: der Teil, den niemand implementiert

Backpressure ist der Mechanismus, mit dem ein langsamer Consumer einem schnellen Producer signalisiert, dass er pausieren soll. Lässt man das weg, puffert man unbegrenzt Daten im Speicher — genau das, was man eigentlich vermeiden wollte.

Wenn du writable.write(chunk) aufrufst, liefert das einen Boolean zurück. false bedeutet, dass der interne Puffer voll ist und du die Quelle pausieren solltest. Diesen Rückgabewert zu ignorieren ist der häufigste Streaming-Bug in produktivem Node.js-Code.

tstypescript
// ❌ Ignores backpressure — the writable's internal buffer grows without bound
readable.on("data", (chunk) => {
  writable.write(chunk); // return value silently discarded
});
 
// ✅ Respects backpressure — pauses the source when the consumer's buffer is full
readable.on("data", (chunk) => {
  const canContinue = writable.write(chunk);
  if (!canContinue) {
    readable.pause();
    writable.once("drain", () => readable.resume());
  }
});

stream.pipeline und .pipe() erledigen das automatisch, wenn Streams über sie verbunden sind. Gefährlich wird es, wenn du manuell aus einem Readable liest — über "data"-Events oder .read() — und selbst in einen Writable schreibst. Genau dort geht das Backpressure verloren.

Asynchrone Generatoren: ein saubereres mentales Modell

Die Streams-API ist mächtig, aber umständlich. Asynchrone Generatoren bieten für dasselbe Konzept eine ergonomischere Schnittstelle. Node.js-Readable-Streams implementieren AsyncIterable nativ, sodass du sie direkt mit for await...of durchlaufen kannst.

tstypescript
async function* transformRows<T, R>(
  source: AsyncIterable<T>,
  transform: (row: T) => R | Promise<R>,
): AsyncGenerator<R> {
  for await (const row of source) {
    yield await transform(row);
  }
}
 
async function* serializeToNDJSON<T>(source: AsyncIterable<T>): AsyncGenerator<string> {
  for await (const item of source) {
    yield JSON.stringify(item) + "\n";
  }
}

Diese lassen sich ohne Umstände kombinieren:

tstypescript
async function streamExport(tenantId: string, res: ServerResponse): Promise<void> {
  const rows = db.queryStream<RawUser>(
    "SELECT id, email, name, created_at FROM users WHERE tenant_id = $1 ORDER BY created_at",
    [tenantId],
  );
 
  const normalized = transformRows(rows, normalizeUser);
  const serialized = serializeToNDJSON(normalized);
 
  try {
    for await (const chunk of serialized) {
      const ok = res.write(chunk);
      if (!ok) {
        // Await drain before pulling more from the generator
        await new Promise<void>((resolve) => res.once("drain", resolve));
      }
    }
  } finally {
    res.end();
  }
}

Die for await...of-Schleife pausiert von selbst, sobald du darin await verwendest — Backpressure gibt es gratis dazu, solange du beim Schreiben wartest (await), wenn der Puffer voll ist.

!

Asynchrone Generatoren propagieren Stream-Fehler nicht automatisch. Umschließe deinen Consumer immer mit try/finally und zerstöre die vorgelagerte Quelle im Fehlerfall — Datenbank-Cursor und Datei-Handles leaken sonst unbemerkt, wenn der Readable mitten in der Iteration aufgegeben wird.

Praxis: einen CSV-Export streamen

Ein vollständiger CSV-Export-Handler zeigt alle Teile zusammen in produktionsreifer Form:

tstypescript
import { Transform, pipeline } from "node:stream";
import { promisify } from "node:util";
 
const pipelineAsync = promisify(pipeline);
 
class CsvTransform extends Transform {
  private headerWritten = false;
 
  constructor(private readonly headers: string[]) {
    super({ objectMode: true });
  }
 
  override _transform(
    row: Record<string, unknown>,
    _encoding: BufferEncoding,
    callback: (err?: Error | null) => void,
  ): void {
    try {
      if (!this.headerWritten) {
        this.push(this.headers.join(",") + "\r\n");
        this.headerWritten = true;
      }
 
      const values = this.headers.map((h) => {
        const raw = String(row[h] ?? "").replace(/"/g, '""');
        return raw.includes(",") || raw.includes('"') || raw.includes("\n")
          ? `"${raw}"`
          : raw;
      });
 
      this.push(values.join(",") + "\r\n");
      callback();
    } catch (err) {
      callback(err as Error);
    }
  }
}
 
export async function handleCsvExport(
  tenantId: string,
  res: NodeJS.WritableStream,
): Promise<void> {
  const queryStream = db.queryStream<Record<string, unknown>>(
    "SELECT id, email, name, created_at FROM users WHERE tenant_id = $1",
    [tenantId],
  );
 
  const csvTransform = new CsvTransform(["id", "email", "name", "created_at"]);
 
  // The database cursor, transform, and HTTP response are all destroyed together on error
  await pipelineAsync(queryStream, csvTransform, res);
}

Der Heap bleibt bei jeder Ergebnisgröße flach. Trennt der Client die Verbindung mitten im Download, propagiert sich der Fehler zurück durch pipeline, und der Datenbank-Cursor wird sofort freigegeben.

Wann man nicht streamen sollte

Streaming bringt echte Komplexität mit sich. Nicht jede große Antwort braucht es.

SzenarioAnsatzGrund
< 10.000 Zeilen, interne APIIn den Speicher ladenEinfacher, vernachlässigbare Heap-Auswirkung
Paginierter UI-EndpunktCursor-basierte PaginierungClient braucht wahlfreien Zugriff pro Seite
Großer Export (CSV, NDJSON)StreamingKonstanter Speicherverbrauch, kein Antwort-Timeout
ETL, Millionen ZeilenStreaming + Batch-SchreibvorgängeSowohl Durchsatz als auch Speicher zählen
Echtzeit-Dashboard-FeedSSE oder WebSocketEin völlig anderes Problem

Die Faustregel: Streame, wenn du die Ergebnisgröße zum Zeitpunkt der Anfrage nicht eingrenzen kannst oder wenn der Client mit dem Konsumieren beginnen soll, bevor der Producer fertig ist. Greife erst nach dem Profiling zu Streaming, wenn es sich als nötig erweist — nicht standardmäßig.

Die wichtigsten Erkenntnisse

  1. stream.pipeline statt .pipe() — korrekte Fehlerpropagierung und Stream-Aufräumung, immer.
  2. Backpressure ist nicht optional — prüfe den Boolean-Rückgabewert von writable.write() und pausiere die Quelle bei false.
  3. Asynchrone Generatoren lassen sich natürlich kombinieren — for await...of macht Backpressure ergonomisch, wenn du Schreibvorgänge innerhalb der Schleife awaitest.
  4. Quellen im Fehlerfall zerstören — Datenbank-Cursor und Datei-Handles leaken still und leise, wenn Readables aufgegeben werden; sowohl pipeline als auch try/finally schützen davor.
  5. Vor dem Streamen messen — eine Pipeline zu einer 200-Zeilen-Abfrage hinzuzufügen bringt Komplexität ohne jeden Nutzen; der richtige Zeitpunkt für Streams ist, wenn die Ergebnisgröße unbegrenzt ist oder große Exporte eine echte Anforderung sind.
Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX