Node.js pipeline 失败收尾:错误已经捕获,为什么三段流都被销毁了
把读取、转换、写入串成 pipeline 后,最外层捕获到了错误,并不表示内部流仍然适合继续使用。pipeline 会协调这组流的失败收尾。另一个同样重要的事实是:已经送到目标的内容,不会因为整体失败而自动消失。
下面构造三个内存流,让转换器遇到 bad 时通过回调报告错误。目标把此前接收的值放入数组,因此能直接检查“任务失败”和“已经发生部分写入”如何同时成立。程序使用默认自动销毁行为,不连接文件、网络或数据库。
在 Node.js v24.19.0 上保存为 demo.mjs,执行 node demo.mjs。使用 node:stream/promises 的 Promise 版 pipeline;assert.rejects 是实验中的失败断言,真实业务可在 await 外用 try/catch记录并安排后续处理。
AI模型生成概念插图:中间站失败后连接的通道关闭,末端托盘仍留着先前完成的块;块数是概念比喻,不是性能或数据统计图。
完整实验程序
import assert from 'node:assert/strict';
import { Readable, Transform, Writable } from 'node:stream';
import { pipeline } from 'node:stream/promises';
function makePipeline(input) {
const received = [];
const source = Readable.from(input, { objectMode: true, highWaterMark: 1 });
const transform = new Transform({
objectMode: true,
transform(value, encoding, callback) {
if (value === 'bad') callback(new Error('invalid record'));
else callback(null, value.toUpperCase());
}
});
const sink = new Writable({
objectMode: true,
write(value, encoding, callback) { received.push(value); callback(); }
});
return { source, transform, sink, received };
}
const failed = makePipeline(['one', 'bad', 'three']);
await assert.rejects(
pipeline(failed.source, failed.transform, failed.sink), /invalid record/
);
console.log('pipeline', 'rejected: invalid record');
const states = [failed.source, failed.transform, failed.sink]
.map(stream => stream.destroyed);
assert.deepEqual(states, [true, true, true]);
assert.deepEqual(failed.received, ['ONE']);
assert.equal(failed.sink.writableFinished, false);
console.log('destroyed', JSON.stringify(states));
console.log('partial output', JSON.stringify(failed.received));
console.log('sink finished', failed.sink.writableFinished);
const fresh = makePipeline(['one', 'two', 'three']);
await pipeline(fresh.source, fresh.transform, fresh.sink);
assert.deepEqual(fresh.received, ['ONE', 'TWO', 'THREE']);
assert.equal(fresh.sink.writableFinished, true);
console.log('fresh output', JSON.stringify(fresh.received));
console.log('fresh finished', fresh.sink.writableFinished);本地实际输出
pipeline rejected: invalid record destroyed [true,true,true] partial output ["ONE"] sink finished false fresh output ["ONE","TWO","THREE"] fresh finished true
关闭链路与撤销输出是两回事
destroyed 的三个 true 依次对应source、transform、sink。转换器用 callback(new Error(...)) 报告失败,pipeline 返回的Promise拒绝;本例随后观察到三个流均已销毁。不要在 catch 中继续向同一个sink写入错误说明,并假定它还能正常接收。
partial output 里保留了 ONE,sink finished 却是 false。第一条数据已经被write回调加入数组,第二条在转换处失败,后续没有正常完成整条流。这份部分结果是程序真实执行的效果,不会因为assert.rejects验证成功就被回滚。
如果目标是生成文件,可由业务决定写临时目标并在整体成功后发布,或保留部分文件用于诊断;如果目标产生外部业务动作,更要事先设计幂等与补偿。pipeline管理的是流的传输和结束,不提供任意输出的事务回滚。
下一次处理创建新的流实例
fresh 再次调用工厂,得到新的源、转换器和目标,把修正后的三条数据全部转换为大写。fresh finished 为 true,与失败分支形成对照。修复输入后应重新建立合适的处理链,而不是把已销毁的实例当作可以按重试按钮复位的对象。
官方规则对已经发出结束或关闭事件的流有例外,不能把“所有场景都再调用一次destroy”当成保证。本文断言只针对这组三个默认设置的流;更换autoDestroy、定制资源释放或特殊流实现后,应核对相应生命周期。
pipeline失败后还可能在部分流上留下监听器,重复利用失败的实例容易让调试更混乱。涉及HTTP响应时尤其要预先安排错误处理:如果socket已被销毁,后来再尝试发送友好错误页可能已经太晚。先确认目标在失败时会处于什么状态,再决定错误应在何处变成用户可见结果。
参考资料
官方资料核验日期:2026-10-02。以上输出来自文中完整程序,断言通过,退出状态为零。


