先记住这个答案
对正常可写流,write 返回 false 表示内部待处理量达到相关阈值,调用方应暂停继续写入,等待 drain 后再生产下一块。当前 chunk 已被接收进写入流程,不能因为 false 就重复提交,否则会重复数据。这个返回值也不是成功落盘或远端确认,写入仍可能随后失败。需要完整处理 error、关闭与结束;如果目标在等待期间被销毁,不能无限等 drain。使用 pipeline 等组合工具时,可以让标准流机制帮助传播背压。
- false 通常是暂停生产的信号
- 当前 chunk 不应因 false 再写一次
- drain 不等于落盘或业务确认
背压限制的是生产者领先消费者的程度
慢磁盘、网络或转换器无法无限接收上游数据,Writable 用内部缓冲和阈值向调用方反馈压力。如果忽略 false 持续写入,Node 不会自动替业务停止所有生产,队列和内存可能继续增长。
等待 drain 后应从下一块继续,而不是重发导致 false 的那块。drain 表示可以恢复写入的时机,并不证明数据已经持久化到磁盘或被远端业务处理,需要这些保证时还要使用相应协议。
用慢速消费者观察每一块只写一次
下面对象模式的阈值设为一,每次写入通过异步回调完成。生产者先写当前值,收到 false 后等待 drain,再取下一个值,最后通过 finished 等待正常结束。
示例消费者固定成功,用于验证节奏和顺序。真实消费者可能出错,等待 drain 时需要同时处理错误或销毁;这里的 once 在 error 事件到来时会拒绝等待,不能把它简化成无条件最终会发生的通知。
import { Writable } from 'node:stream';
import { finished } from 'node:stream/promises';
import { once } from 'node:events';
export async function writeThree() {
const received: string[] = [];
const writeResults: boolean[] = [];
const sink = new Writable({
objectMode: true,
highWaterMark: 1,
write(value: string, _encoding, callback) {
setTimeout(() => { received.push(value); callback(); }, 5);
}
});
for (const value of ['A', 'B', 'C']) {
const ready = sink.write(value);
writeResults.push(ready);
if (!ready) await once(sink, 'drain');
}
const done = finished(sink, { cleanup: true });
sink.end();
await done;
return { writeResults, received };
}测试应确认收到的值恰好是 A、B、C,各出现一次。较小阈值让背压容易观察,但不是通用性能推荐;生产阈值需要结合对象大小、消费者速度与并发数量评估。
最终结束与中途恢复要分别处理
生产者没有更多数据时调用 end,之后应等待最终完成,而不是再用 drain 作为结束证明。drain 是继续写入的信号,finish 或合适的完成工具才用于判断可写侧是否正常完成。
如果下游提前关闭且没有正常 drain,手写等待逻辑需要取消或失败通道。复杂链路优先考虑 pipeline,并测试中途错误、取消和短路退出,避免只在理想情况下控制了速度,却在失败时留下悬挂任务。
容易答错的地方
- 收到 false 就重写相同 chunk
- 正常背压下当前数据已经进入处理流程,再写一次会制造重复。应保存生产游标并等待恢复后推进到下一块,测试时检查内容和次数,而不能只看最后任务有没有结束。
- 把 drain 当作可靠落盘确认
- drain 反映流缓冲的可继续写入状态,具体存储或网络确认有各自语义。需要持久化或远端确认时应使用对应 API 和业务协议,不能从缓冲恢复推导更强的数据可靠性。
面试官还会怎么问?
write 返回 true 就表示不会发生错误吗?
不是。它主要表示当前是否建议继续生产,实际写入仍可能异步失败。应监听错误并处理写回调或整个流的完成结果,不能把布尔返回值当成数据处理成功的最终确认。
为什么 async data 回调里 await 写入仍可能积压?
事件监听器返回的 Promise 不会自动暂停源流,新的 data 仍可能到达并启动更多任务。应使用明确的背压连接或按顺序消费机制,而不是假设给回调加 async 就控制了上游速度。
增大 highWaterMark 能彻底解决背压吗?
不能,它只是允许更多待处理数据,可能降低短时等待频率但增加内存和排队延迟。若消费者长期更慢,积压仍然存在,应控制生产速度、并发和下游吞吐,而不是无限扩大缓冲。
参考资料
示例用于理解所注明的运行环境与边界;延伸学习可结合原文中的更多案例。