13 pipeline():推荐的 Stream 组合方式
1. Promise Pipeline
Node.js 22 中推荐:
import { createReadStream, createWriteStream } from "node:fs";
import { pipeline } from "node:stream/promises";
await pipeline(
createReadStream(sourcePath),
createWriteStream(targetPath),
);
Promise resolve 表示 Pipeline 中的 Stream 已按其完成语义结束;任何阶段失败则 reject。
2. 为什么比 pipe() 好
pipeline() 统一处理:
- 背压;
- 多段错误传播;
- 失败时销毁参与的 Stream;
- 等待完成并返回一个 Promise;
- 避免为每一段重复写
error监听器。
它不能自动删除磁盘半成品,也不能回滚已提交的数据库写入。这些仍属于业务层。
3. 正确清理半成品
import { rm } from "node:fs/promises";
export async function copyFileSafely(
sourcePath: string,
targetPath: string,
): Promise<void> {
try {
await pipeline(
createReadStream(sourcePath),
createWriteStream(targetPath, { flags: "wx" }),
);
} catch (error: unknown) {
await rm(targetPath, { force: true }).catch(() => undefined);
throw error;
}
}
生产代码应记录清理失败,不能永远吞掉;这里的 catch 只为了不覆盖原始错误。
4. 使用 AbortSignal
const controller = new AbortController();
await pipeline(
createReadStream(sourcePath),
createWriteStream(targetPath),
{ signal: controller.signal },
);
取消:
controller.abort();
Promise 会以 AbortError 失败,Pipeline 负责销毁参与 Stream。业务层随后清理半成品并更新任务状态。
5. 超时
const controller = new AbortController();
const timeout = setTimeout(() => {
controller.abort();
}, 30_000);
try {
await pipeline(input, output, {
signal: controller.signal,
});
} finally {
clearTimeout(timeout);
}
Node.js 也提供 AbortSignal 相关辅助能力,可按当前版本选择。重点是超时必须真正传递到工作单元,不能只向客户端返回错误后让文件继续写。
6. Transform 链
await pipeline(
createReadStream(sourcePath),
createGzip(),
createWriteStream(`${sourcePath}.gz`),
);
任一段报错,整个 Promise 失败。
7. finished() 的用途
如果某个库已经为你建立了 Stream 链,只需要等待单个 Stream 的完成,可以使用:
import { finished } from "node:stream/promises";
await finished(output);
但能由你控制全链路时,pipeline() 更能表达“这次传输作为整体成功”。
8. Pipeline 与事务
Pipeline 失败
→ 能停止/销毁 Stream
→ 不能撤销已经写入 MySQL 的记录
应先把数据库任务设为 processing,文件成功完成后再更新为 ready;失败则更新为 failed。若数据库更新失败,还要由补偿任务发现孤立文件。
练习题
- 用 Pipeline 重写上一章的手工
pipe()Promise。 - 中途 Abort,确认两个 Stream 被销毁且半成品被清理。
- 在 gzip Transform 中制造错误,观察 Promise。
- 解释 Pipeline 成功与数据库事务成功为何不是同一件事。