Node.js 写入背压:write 返回 false,为什么这一块不能再发一次

前天 3阅读

把大量记录写进一个慢速目标时,内存一直上涨,常见原因是循环只顾调用 write,完全没有理会它返回的布尔值。另一种修复也会出错:看见 false,就把刚才那条记录再写一次。这样队列不但没有减压,还会产生重复数据。

对正常可写、尚未结束的 Writable,false 表示这次写入已经进入处理流程,同时要求生产者先暂停后续写入。恢复信号是 drain。把“这条是否已接收”和“下一条能否继续送”拆开,才不会把流量控制误写成重试机制。

把慢消费者做成可观察的实验

保存为 demo.mjs,执行 node demo.mjs。本文在 Node.js v24.19.0 实跑。目标只把字符放入内存数组,并通过 setImmediate 延后完成写入,不依赖磁盘速度或网络条件。对象模式的阈值设为二,计量单位是对象数量。

Node.js 写入背压:write 返回 false,为什么这一块不能再发一次

AI概念插图:用抽象图形表现本文讨论的关系,并非软件截图或真实运行结果。

import assert from 'node:assert/strict';
import { Writable } from 'node:stream';
import { finished } from 'node:stream/promises';
import { once } from 'node:events';

const received = [];
const sink = new Writable({
  objectMode: true,
  highWaterMark: 2,
  write(value, encoding, callback) {
    setImmediate(() => { received.push(value); callback(); });
  }
});
const done = finished(sink, { cleanup: true });
done.catch(() => {});
for (const value of ['A', 'B', 'C']) {
  const canContinue = sink.write(value);
  console.log('write', value, canContinue);
  if (!canContinue) {
    await once(sink, 'drain');
    console.log('drain');
  }
}
sink.end();
await done;
assert.deepEqual(received, ['A', 'B', 'C']);
console.log('finished', received.join(','));

const burst = new Writable({
  objectMode: true, highWaterMark: 2,
  write(value, encoding, callback) { setImmediate(callback); }
});
const burstDone = finished(burst, { cleanup: true });
for (let n = 0; n < 5; n++) burst.write(n);
assert.equal(burst.writableLength, 5);
console.log('ignored-pressure', burst.writableLength);
burst.end();
await burstDone;

const broken = new Writable({
  objectMode: true, highWaterMark: 1,
  write(value, encoding, callback) {
    setImmediate(() => callback(new Error('demo-failure')));
  }
});
const failed = finished(broken, { cleanup: true }).catch(e => e.message);
assert.equal(broken.write('X'), false);
await assert.rejects(once(broken, 'drain'), /demo-failure/);
assert.equal(await failed, 'demo-failure');
console.log('error rejected drain wait');

false 的那一条仍然只写一次

前三次写入依次输出 A true、B false、C true,中间出现一次 drain。第二次调用把 B 交给目标后,程序停在等待处;收到恢复通知,才继续发送 C。最后 finished 后的内容是 A,B,C,断言同时检查了顺序与没有重复。

等待应紧跟在返回 false 的调用之后建立。示例里的 once 也会在目标发出 error 时拒绝等待,末尾的故障目标实际验证了这个分支。不要把 drain 当成必达通知:异常、销毁或其他终止条件都需要另行处理,不能无限等待一个已停止的目标。

done 在写入前就注册了 finished 观察。提前附加 catch 只是防止拒绝暂时无人处理,原来的 done 仍保留失败状态,后面的 await 仍会抛出。流错误不会因为加了这个空处理器就被当作成功。

日志里的布尔值不是每条记录的成功回执。要确认记录是否真正处理,应观察目标的完成回调或任务结束结果;如果业务还需要逐条确认,则把记录标识和确认协议单独设计,不能拿流的缓存状态替代它。这个区别在慢速批量导出时尤其容易被忽略。

阈值、恢复和完成各自回答不同问题

第二段故意连续写五个对象,不遵守返回值。同步读取的 writableLength 是五,超过设置的二。这个小反例说明 highWaterMark 是发出背压提示的阈值,不是替你拒收更多数据的硬内存上限;无限忽略提示会让缓存不断增长。

drain 表示可以恢复写入,不表示整项任务已经完成。所有输入送出后还要调用 end,并等待 finished。本例能证明内存目标已按约定处理输入;如果换成文件或远端服务,是否持久化、是否得到业务确认,还要看对应接口的保证。

这里选择对象模式是为了让计数容易复核。普通字节流的阈值通常按字节计算,不能直接把同一个数字理解成记录条数。一个大对象也可能携带很多数据,对象数量受控并不自动等于总体内存受控。

接入自己的数据源时,先用三条不同标识的记录复现暂停和恢复,再测试目标报错。若是现成流之间的传输,可评估 pipeline 来协调背压和结束;本篇手动循环只验证这个受控目标,不作为覆盖所有关闭竞态的通用泵送器。

参考资料

文章版权声明:除非注明,否则均为云鹊BLOG原创文章,转载或复制请以超链接形式并注明出处。