バックプレッシャー、pipe()、pipeline()ユーティリティ
メモリ膨張を診断し、pipeline()でストリームを安全に接続して、エラーを伝播させバックプレッシャーに対応します。
「バックプレッシャー、pipe()、pipeline()ユーティリティ」はCoddyKit上の無料Node.js Backend Development Bootcampレッスンです。 これはレッスン3/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはNode.js Backend Development Bootcamp学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Node.js Backend Development Bootcampコースには全4レッスンが含まれています。
このレッスンの一部はまだ翻訳されておらず、英語で表示されています。
Why Memory Bloats in Streams
A Node.js Readable stream can produce data faster than a Writable can consume it. If you never tell the producer to slow down, unconsumed chunks pile up in an internal buffer and your process memory grows until the GC can't keep up.
- A slow disk, slow network socket, or slow database write is the typical consumer.
- A fast file read or HTTP upload is the typical producer.
The mechanism that makes the producer wait for the consumer is called backpressure. Misusing streams almost always means backpressure was ignored.
The Naive (Broken) Copy
Here is the classic memory bug. We listen for data and call dst.write() for every chunk, ignoring its return value.
If dst is slower than src, the unwritten chunks queue up inside dst's buffer with no upper bound. For a multi-gigabyte file this can exhaust RAM.
const fs = require('fs');
const src = fs.createReadStream('big.bin');
const dst = fs.createWriteStream('copy.bin');
// BUG: return value of write() is ignored, so backpressure is never honored
src.on('data', (chunk) => {
dst.write(chunk);
});
src.on('end', () => dst.end());What write() Actually Returns
writable.write(chunk) returns a boolean:
true— the internal buffer is belowhighWaterMark; keep writing.false— the buffer is full; you should stop writing and wait for the'drain'event before sending more.
Honoring this return value is the manual way to apply backpressure. The producer must pause until the consumer signals it has drained.
Manual Backpressure with pause/resume
Done by hand, backpressure means: when write() returns false, pause() the source; when the destination emits 'drain', resume() it.
This works but is verbose and easy to get wrong — you also have to wire up error and end handling for both streams.
const fs = require('fs');
const src = fs.createReadStream('big.bin');
const dst = fs.createWriteStream('copy.bin');
src.on('data', (chunk) => {
const ok = dst.write(chunk);
if (!ok) {
src.pause(); // stop reading until the buffer drains
dst.once('drain', () => src.resume());
}
});
src.on('end', () => dst.end());pipe() Does This For You
readable.pipe(writable) wires up the same pause/resume/drain dance automatically and honors backpressure out of the box.
It returns the destination stream, so you can chain through transforms:
src.pipe(gzip).pipe(dst)
For most simple copies, pipe() is far better than the manual loop above.
const fs = require('fs');
const zlib = require('zlib');
const src = fs.createReadStream('big.bin');
const gzip = zlib.createGzip();
const dst = fs.createWriteStream('big.bin.gz');
// pipe() handles backpressure across all three streams
src.pipe(gzip).pipe(dst);The Hidden Flaw in pipe()
pipe() handles backpressure, but it does not forward errors. If gzip or dst emits 'error', the source is not destroyed automatically.
- The upstream stream keeps its file descriptor open — a resource leak.
- An unhandled
'error'event throws and can crash the process.
To use pipe() safely you must attach an error handler to every stream and manually destroy the others. That boilerplate is exactly what pipeline() removes.
Enter stream.pipeline()
stream.pipeline() connects a series of streams, propagates backpressure, forwards errors, and destroys every stream in the chain when any of them fails or finishes.
It takes the streams in order followed by a callback that fires once with an error (or null on success):
const { pipeline } = require('stream');
const fs = require('fs');
const zlib = require('zlib');
pipeline(
fs.createReadStream('big.bin'),
zlib.createGzip(),
fs.createWriteStream('big.bin.gz'),
(err) => {
if (err) {
console.error('Pipeline failed:', err.message);
} else {
console.log('Pipeline succeeded');
}
}
);The Promise-Based pipeline()
In modern code use the promise version from stream/promises. It resolves on success and rejects on failure, so a single try/catch covers the whole chain and cleanup.
This is the recommended way to wire streams in async backend handlers.
const { pipeline } = require('stream/promises');
const fs = require('fs');
const zlib = require('zlib');
async function compress() {
try {
await pipeline(
fs.createReadStream('big.bin'),
zlib.createGzip(),
fs.createWriteStream('big.bin.gz')
);
console.log('done');
} catch (err) {
console.error('failed:', err.message);
}
}
compress();A Runnable In-Memory Pipeline
You don't need files to see pipeline() work. Readable.from() turns any iterable into a stream, and a Transform can uppercase each chunk. The whole thing runs standalone.
Notice how errors from any stage would reject the awaited pipeline().
const { Readable, Transform } = require('stream');
const { pipeline } = require('stream/promises');
const source = Readable.from(['hello ', 'stream ', 'world']);
const upper = new Transform({
transform(chunk, _enc, cb) {
cb(null, chunk.toString().toUpperCase());
}
});
const chunks = [];
const sink = new Transform({
transform(chunk, _enc, cb) {
chunks.push(chunk.toString());
cb();
}
});
(async () => {
await pipeline(source, upper, sink);
console.log(chunks.join(''));
})();highWaterMark: Tuning the Buffer
Each stream has a highWaterMark (default 16 KB for byte streams, 16 objects for object-mode). It is the threshold at which write() returns false and reads pause.
- A larger highWaterMark increases throughput but uses more memory per stream.
- A smaller one applies backpressure sooner, capping memory more tightly.
It is a buffering threshold, not a hard limit — but it is the lever that controls how aggressively backpressure kicks in.
const fs = require('fs');
// Pause reads after only 64 KB is buffered downstream
const src = fs.createReadStream('big.bin', { highWaterMark: 64 * 1024 });
const dst = fs.createWriteStream('copy.bin', { highWaterMark: 64 * 1024 });
src.pipe(dst);pipeline() in an HTTP Handler
A common backend mistake is buffering an entire upload or download into memory before responding. Streaming the response body with pipeline() keeps memory flat and tears everything down if the client disconnects.
Because the HTTP response is a Writable, backpressure from a slow client automatically throttles the file read.
const http = require('http');
const fs = require('fs');
const { pipeline } = require('stream');
http.createServer((req, res) => {
pipeline(
fs.createReadStream('big.bin'),
res,
(err) => {
if (err) {
console.error('stream error:', err.message);
res.destroy();
}
}
);
}).listen(3000);Quick Check
You are streaming a file to a slow client through a gzip transform. Which approach safely honors backpressure AND cleans up every stream if the client disconnects mid-transfer?
Recap
Key takeaways for wiring streams safely:
- Backpressure stops a fast producer from overwhelming a slow consumer; ignoring
write()'s boolean return is the root cause of stream memory bloat. pipe()handles backpressure but not error forwarding or cleanup — a leaked-FD trap.stream.pipeline()(callback or thestream/promisesversion) propagates backpressure, forwards errors, and destroys all streams in the chain.highWaterMarktunes how soon backpressure engages, trading memory for throughput.- In HTTP handlers, stream with
pipeline()instead of buffering full payloads.
よくある質問
「バックプレッシャー、pipe()、pipeline()ユーティリティ」レッスンは無料ですか?
はい。「バックプレッシャー、pipe()、pipeline()ユーティリティ」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Node.js Backend Development Bootcampコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Node.js Backend Development Bootcampコースには全4レッスンが含まれています。
「バックプレッシャー、pipe()、pipeline()ユーティリティ」で何を学びますか?
メモリ膨張を診断し、pipeline()でストリームを安全に接続して、エラーを伝播させバックプレッシャーに対応します。 ブラウザで直接実行するハンズオンコードでNode.js Backend Development Bootcampを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
Node.js Backend Development Bootcampを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのNode.js Backend Development Bootcampは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン3/4です。
「バックプレッシャー、pipe()、pipeline()ユーティリティ」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このNode.js Backend Development Bootcampレッスンでコードを書いて実行できますか?
はい。すべてのNode.js Backend Development Bootcampレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- Readable、Writable、Duplex、Transformストリームの内部
- _transformと_flushによるカスタムTransformストリームの実装
- バックプレッシャー、pipe()、pipeline()ユーティリティ
- ストリーム上の非同期イテレーターとfor-await-of