Skip to content

perf(streams): stream transforms compress and decompress synchronously on the main thread #554

Description

@derodero24

Problem

The stream APIs do all compression work synchronously on the JS thread:

  • Every ctx.transform() / flush() / finish() is a synchronous napi call, made inline from the Web TransformStream callbacks in streams.js (e.g. createBrotliCompressStream) and from the Node Transform callbacks in node.js (e.g. createZstdCompressTransform).
  • Each chunk blocks the event loop for its full compression time.
  • When the source is in memory, the transform callbacks complete synchronously, so the stream machinery moves on to the next chunk through microtasks and process.nextTick. Timers and I/O do not run until the whole stream has been processed.

node:zlib streams run each chunk on the libuv thread pool, so the event loop keeps turning.

The README does not mention this. Its "Choosing an API mode" table (README.md#L200-L207) recommends Async "when event loop must stay free (servers, UIs)" and Streaming for "unknown/unbounded data size". Users of the streaming APIs in servers, including @derodero24/comprs-middleware, can reasonably assume that streaming does not block either.

Reproduction

Linux x64, 4 vCPU, Node 22.22.0, comprs 2.0.2 release build. The input is 16 MB of random words from a 4,096-word vocabulary. A 1 ms setInterval counts how often the event loop got to run timers:

const { createBrotliCompressTransform } = require('@derodero24/comprs/node');
const zlib = require('zlib');
const { Readable, Writable, pipeline } = require('stream');

const vocab = Array.from({ length: 4096 }, (_, i) => i.toString(36).padStart(3, 'x') + 'etaoinshrd'.slice(0, i % 7));
const idx = require('crypto').randomBytes(8 << 20);
const parts = [];
for (let i = 0, n = 0; n < 16 << 20; i += 2) { const w = vocab[idx.readUInt16LE(i % idx.length) & 4095] + ' '; parts.push(w); n += w.length; }
const input = Buffer.from(parts.join('')).subarray(0, 16 << 20);

function run(label, transform) {
  return new Promise((resolve, reject) => {
    let ticks = 0, last = performance.now(), maxGap = 0;
    const timer = setInterval(() => { const now = performance.now(); maxGap = Math.max(maxGap, now - last); last = now; ticks++; }, 1);
    const t0 = performance.now();
    pipeline(Readable.from([input]), transform, new Writable({ write(_c, _e, cb) { cb(); } }), (err) => {
      clearInterval(timer);
      maxGap = Math.max(maxGap, performance.now() - last);
      if (err) return reject(err);
      console.log(label, Math.round(performance.now() - t0), 'ms,', ticks, 'ticks, longest gap', Math.round(maxGap), 'ms');
      resolve();
    });
  });
}

(async () => {
  await run('comprs', createBrotliCompressTransform(9));
  await run('node:zlib', zlib.createBrotliCompress({ params: { [zlib.constants.BROTLI_PARAM_QUALITY]: 9 } }));
})();
Brotli quality 9, 16 MB Wall time Timer ticks Longest event-loop gap
comprs Node transform, 1 chunk of 16 MB 5,711 ms 0 5,711 ms
node:zlib, 1 chunk of 16 MB 4,679 ms 3,982 15 ms
comprs Node transform, 256 × 64 KiB from memory 6,532 ms 0 6,532 ms
node:zlib, 256 × 64 KiB from memory 5,798 ms 4,982 25 ms
comprs Node transform, 256 × 64 KiB from fs.createReadStream 5,766 ms 103 134 ms
node:zlib, 256 × 64 KiB from fs.createReadStream 5,305 ms 4,622 12 ms
comprs Web stream (createBrotliCompressStream(9)), 256 × 64 KiB from memory 5,558 ms 0 5,558 ms

The LZ4 decompression stream is the extreme case. It buffers all input and decodes everything in flush() (tracked separately in #535). For 100 MB of output, that flush() alone blocked for 124-142 ms in three runs.

Proposed fix

  • Async methods. Add transformAsync(chunk) / flushAsync() / finishAsync() to the context classes, implemented as AsyncTasks:
    • Keep the codec state in an Arc<Mutex<...>> (or move it into the task and back), so compute() runs on the thread pool.
    • The JS wrappers already serialize calls per stream, so the lock is uncontended.
    • Reject calls made while another call is still in flight, instead of blocking on the lock.
  • Wrappers. Use the async methods from the Node transform / flush callbacks (call callback when the promise settles) and from the Web transform / flush callbacks (return the promise).
  • Small chunks. Thread-pool dispatch costs tens of microseconds per call. The wrappers could keep the synchronous path for chunks below a threshold, e.g. 16-64 KiB, as long as the total synchronous work per tick stays bounded.
  • Input copies. The chunk-input copy question is the same as for the one-shot *Async functions (perf(core): *Async functions copy their input on the main thread #548).
  • Until then, document the behaviour.
    • State in the README API-mode table and in the stream API docs that streaming runs on the calling thread and blocks the event loop for each chunk.
    • Recommend *Async or a worker thread for CPU-heavy settings such as brotli quality ≥ 9.

Related

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or requestperformancePerformance benchmarks and optimization

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions