0Pricing
Node.js Backend Development Bootcamp · 课时

使用 _transform 和 _flush 实现自定义转换流

构建可复用的转换流,在数据流经时修改、筛选并聚合数据块。

使用 _transform 和 _flush 实现自定义转换流 是 CoddyKit 上的免费 Node.js Backend Development Bootcamp 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Node.js Backend Development Bootcamp 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Node.js Backend Development Bootcamp 课程共包含 4 节课。

本课时的部分内容尚未翻译,以英文显示。

Why Transform Streams?

A Transform stream is both readable and writable: it consumes input chunks, processes them, and pushes output chunks. It is the right tool whenever data must be mutated as it flows rather than buffered fully in memory.

  • Writable side accepts data via write() / pipe() from upstream.
  • Readable side emits processed data that downstream consumers read.

Typical backend uses: uppercasing/normalizing a payload, gzip-style encoding, CSV-to-JSON conversion, redacting secrets in a log pipeline, or counting bytes — all without loading the whole file or HTTP body into RAM.

The Two Hooks: _transform and _flush

A custom Transform stream is defined by implementing two internal methods. Node calls them for you — you never call them directly.

  • _transform(chunk, encoding, callback) — invoked once per incoming chunk. Do your work, push() any output, then signal completion with callback().
  • _flush(callback) — invoked once, after the last chunk, just before the stream ends. Use it to emit any trailing/aggregated data.

The leading underscore marks them as the framework-facing implementation. Consumers still use the public write, read, and pipe API.

A Minimal Uppercase Transform

The classic starting point: extend the Transform class and override _transform. Each chunk is a Buffer (unless object mode), so convert to a string, transform it, and push the result.

Calling callback() with no error tells Node this chunk is fully processed and it may deliver the next one. Passing the value as the second arg to callback is shorthand for push + callback().

const { Transform } = require('node:stream');

class Upper extends Transform {
  _transform(chunk, encoding, callback) {
    const out = chunk.toString().toUpperCase();
    callback(null, out); // shorthand for this.push(out); callback();
  }
}

const up = new Upper();
up.on('data', (d) => process.stdout.write(d));
up.write('hello ');
up.write('streams\n');
up.end();

callback() Is a Contract

The callback in _transform is mandatory. Until you call it, Node assumes the chunk is still in progress and will not hand you the next one. This is how backpressure flows through your transform.

  • callback() — success, ready for next chunk.
  • callback(err) — emits an 'error' event and destroys the stream.
  • callback(null, data) — pushes data and signals success.

Forgetting to call callback is the #1 bug: the pipeline silently stalls forever with no error.

push() Multiple Times Per Chunk

One input chunk can produce zero, one, or many output chunks. Call this.push() as many times as needed before invoking callback(). This is what makes splitting (e.g. line-by-line) possible.

Below, a single write containing several lines is fanned out into one push per line.

const { Transform } = require('node:stream');

class LineSplitter extends Transform {
  _transform(chunk, encoding, callback) {
    const lines = chunk.toString().split('\n');
    for (const line of lines) {
      if (line.length) this.push(line + ' <<\n');
    }
    callback();
  }
}

const s = new LineSplitter();
s.on('data', (d) => process.stdout.write(d));
s.end('alpha\nbeta\ngamma\n');

Filtering: Drop Chunks by Pushing Nothing

To filter, simply decide not to push. If a chunk should be discarded, call callback() without pushing anything — the data never reaches the readable side.

This pattern is ideal for redacting or dropping records mid-pipeline, e.g. removing log lines that contain a secret token.

const { Transform } = require('node:stream');

class DropSecrets extends Transform {
  _transform(chunk, encoding, callback) {
    const line = chunk.toString();
    if (line.includes('SECRET')) {
      return callback(); // filtered out, nothing pushed
    }
    callback(null, line);
  }
}

const f = new DropSecrets();
f.on('data', (d) => process.stdout.write(d));
f.write('ok line 1\n');
f.write('this has a SECRET token\n');
f.write('ok line 2\n');
f.end();

Object Mode for Structured Records

By default chunks are Buffer/string. Set objectMode: true to push and receive JavaScript objects instead — essential for record-oriented pipelines (JSON rows, DB results, parsed events).

  • writableObjectMode / readableObjectMode can be set independently if input and output types differ.
  • In object mode, each push emits exactly one object regardless of size.
const { Transform } = require('node:stream');

class AddTax extends Transform {
  constructor() { super({ objectMode: true }); }
  _transform(order, encoding, callback) {
    callback(null, { ...order, total: order.price * 1.2 });
  }
}

const t = new AddTax();
t.on('data', (o) => console.log(o));
t.write({ id: 1, price: 100 });
t.write({ id: 2, price: 250 });
t.end();

_flush: Emit Trailing/Aggregated Data

_flush(callback) runs once after the final chunk, before 'end'. It is where you push anything you have been accumulating — a running total, a buffered partial line, or a closing delimiter.

You can push inside _flush just like in _transform. You must call its callback() so the stream can finish.

const { Transform } = require('node:stream');

class Summer extends Transform {
  constructor() { super({ objectMode: true }); this.sum = 0; }
  _transform(num, encoding, callback) {
    this.sum += num;
    callback(); // aggregate, emit nothing yet
  }
  _flush(callback) {
    this.push({ total: this.sum }); // emit once at the end
    callback();
  }
}

const agg = new Summer();
agg.on('data', (o) => console.log(o));
[10, 20, 30, 40].forEach((n) => agg.write(n));
agg.end();

Buffering Partial Lines Across Chunk Boundaries

Chunks do not align with logical records. A line may be split across two chunks, so a robust line-parser keeps a leftover buffer between calls and flushes the remainder in _flush.

This combine-in-_transform, drain-in-_flush pattern is the backbone of real CSV/NDJSON parsers.

const { Transform } = require('node:stream');

class LineParser extends Transform {
  constructor() { super({ readableObjectMode: true }); this.tail = ''; }
  _transform(chunk, encoding, callback) {
    const data = this.tail + chunk.toString();
    const parts = data.split('\n');
    this.tail = parts.pop(); // keep incomplete last segment
    for (const line of parts) this.push(line);
    callback();
  }
  _flush(callback) {
    if (this.tail) this.push(this.tail);
    callback();
  }
}

const p = new LineParser();
p.on('data', (l) => console.log('LINE:', l));
p.write('he');
p.write('llo\nwor');
p.write('ld\nlast');
p.end();

The Functional Shorthand: stream.Transform options

You don't always need a class. The Transform constructor accepts transform and flush functions directly — handy for small, one-off transforms.

Inside these functions, this is still the stream, so this.push() works. The class form is preferred when you want reusable, named, instantiable components; the inline form is great for quick glue.

const { Transform } = require('node:stream');

const csvToRows = new Transform({
  readableObjectMode: true,
  transform(chunk, encoding, callback) {
    for (const line of chunk.toString().trim().split('\n')) {
      const [id, name] = line.split(',');
      this.push({ id: Number(id), name });
    }
    callback();
  },
});

csvToRows.on('data', (row) => console.log(row));
csvToRows.end('1,Ada\n2,Linus\n3,Grace');

Composing in a Pipeline

Transform streams shine when chained. Use stream.pipeline() (callback or promise form) instead of raw .pipe() — it propagates errors and destroys every stream on failure, preventing leaks.

Here a source feeds an uppercaser, then a suffix-adder, then stdout. Each transform stays small and reusable.

const { Transform, Readable, pipeline } = require('node:stream');

const make = (fn) => new Transform({
  transform(chunk, enc, cb) { cb(null, fn(chunk.toString())); },
});

const upper = make((s) => s.toUpperCase());
const bang = make((s) => s + '!\n');

pipeline(
  Readable.from(['log a\n', 'log b\n']),
  upper,
  bang,
  process.stdout,
  (err) => { if (err) console.error('failed', err); else console.error('done'); }
);

Quick Check: Where Do Trailing Aggregates Go?

You are building a Transform that counts total bytes seen and must emit a single summary object after all input is processed. Which method should push that summary?

Recap: Transform Stream Essentials

You can now build reusable Transform streams that mutate, filter, and aggregate flowing data:

  • _transform(chunk, enc, cb) — process each chunk; push zero or more outputs; always call cb() (or cb(null, data)).
  • _flush(cb) — runs once at the end to emit trailing or aggregated data; must call cb().
  • Filter by pushing nothing; split by pushing many times; aggregate by accumulating state and flushing.
  • objectMode (and the independent readable/writable variants) carries structured records.
  • Buffer partial records in _transform and drain them in _flush.
  • Compose with stream.pipeline() for safe error handling and cleanup.

Never forget the callback — a missing cb() silently stalls the entire pipeline.

常见问题解答

「使用 _transform 和 _flush 实现自定义转换流」课时是免费的吗?

是的 — 「使用 _transform 和 _flush 实现自定义转换流」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Node.js Backend Development Bootcamp 课程的其余内容,请升级到 CoddyKit PRO。 Node.js Backend Development Bootcamp 课程共包含 4 节课。

「使用 _transform 和 _flush 实现自定义转换流」这节课中我会学到什么?

构建可复用的转换流,在数据流经时修改、筛选并聚合数据块。 你通过在浏览器中直接运行的动手代码来练习 Node.js Backend Development Bootcamp,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Node.js Backend Development Bootcamp 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Node.js Backend Development Bootcamp 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。

「使用 _transform 和 _flush 实现自定义转换流」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Node.js Backend Development Bootcamp 课中编写并运行代码吗?

能。每节 Node.js Backend Development Bootcamp 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 可读、可写、双工与转换流内部机制
  2. 使用 _transform 和 _flush 实现自定义转换流
  3. 背压、pipe() 与 pipeline() 工具
  4. 异步迭代器与通过 for-await-of 遍历流
← 返回 Node.js Backend Development Bootcamp