本页目录

流与背压:让快生产者等待慢消费者

通过慢写入、转换管道和取消实验理解流的容量反馈,分清结束、关闭与失败,避免把流重新堆成整份内存。

L2 · 能交付约 15 分钟阅读含示例、练习与验收

建议先读:Buffer 与编码:字符串之外的字节世界EventEmitter:同步通知、错误通道与监听器寿命

本页内容

目标与前置#

团队助手要处理的文档可能比可用内存大,也可能来自速度不稳定的网络。一次 readFile 把所有内容放进内存,写起来简单,却把文件规模直接变成内存需求。流提供另一种契约:逐步生产数据,逐步消费数据,并让消费速度反过来限制生产。前置是 Buffer 的字节边界和 EventEmitter 的错误、监听生命周期。

学完以后,你应当能解释 write 返回 false 的含义,使用 pipeline 连接读取、转换和写入,设置明确的容量阈值,并让取消与失败结束整条处理链。我们不会把每个事件都手动组合成生产级框架,而是先通过底层实验看到背压,再使用现代稳定 API 管理共同生命周期。

流首先是协议,其次才是文件工具#

Readable 产生数据,Writable 接收数据,Duplex 同时拥有可读和可写两侧,Transform 则是一种把输入转换为输出的双向流。来源可以是文件、网络或异步生成器,目标可以是磁盘、压缩器或另一个连接。理解这套接口以后,处理逻辑就不必依赖数据最初来自哪里。

流的价值不只是“分成许多块”。如果生产者不断推送,而消费者把块全部保存在数组中,最后再拼成一个大 Buffer,峰值内存仍随总输入增长。真正的流式处理要求每个阶段尽量在处理后释放历史数据,并把不能继续接受的信号传回上游。只用了 data 事件并不能证明已经实现流式内存控制。

数据块也不是业务记录。一个块可能含半行、几行或半个中文字符。字节到文本的解码状态、文本到记录的分隔状态,需要各自维护。转换器不应该假设一次 transform 就对应一个完整文档,否则在本机小样本正常、真实上传分块改变后就会失败。

背压:接受当前块,但请暂时停止继续写#

writable.write(chunk, encoding, callback) 把数据交给可写流,返回布尔值。false 通常表示内部缓冲达到阈值,调用者应该停止继续提交,等待 drain;它不表示当前块被拒绝,因此不能把同一块再写一次,否则会重复数据。回调表示该块处理完成或失败,返回值表达容量反馈,两者承担不同职责。

highWaterMark 是内部缓冲达到何种程度时发出反馈的阈值,不是整个进程的硬内存上限。普通字节流通常按字节衡量,对象模式按对象个数衡量;解码字符串的某些路径还涉及代码单元计数。一个对象可能包含很大数组,因此对象数量受限不等于字节占用受限。

默认阈值会受到流类型、Node 版本与平台影响,Node 22 起通用默认值也有变化。本教材在需要观察容量的地方显式设置阈值,不背诵一个适用于所有流的固定数字。调整阈值会改变批量与缓冲权衡,但不会让慢磁盘突然变快;它只是让系统在等待时存放多少待处理数据。

end() 表示不再提交新数据,随后等待可写侧处理完已经提交的内容。finish 表示可写处理完成,close 表示底层资源关闭,error 表示失败,它们不是可以互换的“结束”。可读侧的 end 则表示数据已经读完。对于文件和网络,读完、写完、关闭与业务成功之间还存在额外约定。

实验一:亲手遵守 write 的反馈#

环境 Node 22.22.0,无依赖。保存为 stream-drain.mjs。它创建一个故意缓慢的可写流,不写磁盘,只统计自己接收的字节。

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

let total = 0;
const sink = new Writable({
  highWaterMark: 4,
  write(chunk, encoding, callback) {
    setTimeout(() => { total += chunk.length; callback(); }, 2);
  }
});
let pauses = 0;
for (let index = 0; index < 5; index += 1) {
  // false 仍已接受本块,只暂停后续提交。
  if (!sink.write(Buffer.from('abcd'))) {
    pauses += 1;
    await once(sink, 'drain');
  }
}
sink.end();
await finished(sink, { cleanup: true });
assert.equal(total, 20);
assert.equal(pauses, 5);
console.log(`字节=${total}, 暂停=${pauses}`);

执行 node stream-drain.mjs,预期 字节=20, 暂停=5。每块四字节恰好达到显式阈值,因此本例每次都暂停;慢写入回调完成后,流有机会发出 drain,再提交下一块。最后 end 并等待 finished,证明我们不仅提交了数据,还等待可写侧真正结束。

构造器的 write 方法是实现者接口,callback 必须在本块处理完成时调用一次,失败则传入错误。不能在内部抛一个异步异常,希望流自动捕获;也不能忘记回调,否则队列会一直等这块。finished 返回 Promise,完成时兑现,失败时拒绝;cleanup 选项让它清理自身事件监听,避免观察器残留。

这个例子只演示一个受控可写流。手工连接多个真实资源时,你还必须传播读错误、写错误、关闭和取消,并避免漏掉任何一方。因此日常组合优先选择 pipeline,而不是复制一串只处理正常 data 与 end 的事件监听。

pipeline 管理的是整条处理关系#

node:stream/promisespipeline(source, transform, destination, options) 返回 Promise,成功时表示管道按约定完成,失败时拒绝并对相关流进行销毁处理。最后的 options 可以包含 AbortSignal。它会协调常见背压与错误传播,让你把注意力放在每个阶段的转换规则上。

pipeline 不会让任意转换变成低内存。一个 Transform 若在内部保存所有历史块,仍然会占用全部数据;一个写入目标若只是把数据放进另外一个无限队列,也只是把积压挪了位置。端到端容量需要每个阶段都遵守边界,包括自定义队列与外部客户端。

管道出错后,相关流通常不能作为全新任务重新使用。为下一次运行创建新的来源、转换器和目标,避免复用已销毁对象或残留监听器。错误也不会自动撤销已经写出的前半段文件;如果业务要求成功后才可见,应该写入临时目标,确认完成后再发布,并在失败时清理自己创建的临时资源。

实验二:转换阶段限制总字节并传播失败#

保存为 stream-pipeline.mjs。每次运行创建全新的流,成功路径验证字节数,失败路径验证上限错误与资源销毁。

stream-pipeline.mjs
import { Readable, Transform, Writable } from 'node:stream';
import { pipeline } from 'node:stream/promises';
import assert from 'node:assert/strict';

async function run(limit) {
  const source = Readable.from([Buffer.from('abcd'), Buffer.from('efgh')], { objectMode: false });
  let seen = 0;
  let saved = 0;
  const guard = new Transform({
    transform(chunk, encoding, callback) {
      seen += chunk.length;
      if (seen > limit) { callback(new RangeError('文档字节超限')); return; }
      callback(null, chunk);
    }
  });
  const sink = new Writable({
    write(chunk, encoding, callback) { saved += chunk.length; callback(); }
  });
  try {
    await pipeline(source, guard, sink);
    return saved;
  } catch (error) {
    assert.equal(guard.destroyed, true);
    assert.equal(sink.destroyed, true);
    throw error;
  }
}
assert.equal(await run(8), 8);
await assert.rejects(run(5), { name: 'RangeError', message: '文档字节超限' });
console.log('成功管道与超限销毁通过');

执行 node stream-pipeline.mjs,预期通过。Readable.from 把可迭代输入适配成流,本例显式关闭对象模式;转换器累计的是实际字节数,不信任来源声明;失败时通过 callback(error) 报告给流协议,pipeline 再把失败传给调用方。成功时 callback(null, chunk) 继续传递同一块,不需要复制。

五字节上限的失败路径可能已经让前面的部分内容经过目标,因此“管道失败”不等于“目标从未看到任何字节”。这个区别对上传和导出尤其重要。需要整体发布的文件要分离临时写入与最终可见状态;需要数据库写入的记录则考虑事务或可恢复的批次协议,不能靠流销毁实现回滚。

取消是有方向的资源请求#

AbortSignal 传入 pipeline 后,取消会使管道以取消错误结束,并销毁相关流。可是异步生成器内部若正在等待一个与信号无关的长任务,仅销毁外层适配器不一定能立刻中止那个等待。生成器也应接收信号,并把它传给支持取消的底层操作。取消需要贯穿整条调用链。

取消与成功可能接近同时发生,所以清理应可重复调用,状态转换应明确。调用方收到 AbortError 通常可以把任务标记为用户取消,而不是普通服务器失败;但已经发生的外部副作用仍要按业务规则处理。不要在收到取消后直接删除任意目标,除非确认目标属于本次尚未提交的操作。

网络流还有额外边界。把 HTTP 请求或响应直接交给 pipeline,错误时可能销毁连接,导致你无法再发送一个格式完整的错误响应。因此读请求体时,要提前决定超限后是排空并返回错误,还是关闭连接;输出已开始后发生失败,也不能再修改已经发送的状态码。后续网络与服务章节会把这些限制放进具体路由。

练习:取消一个对象流任务#

使用异步生成器逐项产生任务编号,每项之间等待五毫秒。目标接收第三项时取消。要求管道以 AbortError 拒绝,生成器的 finally 必须执行,目标数量停在三。提示:把 pipeline 提供的 signal 传给 Promise 版 setTimeout,避免生成器在取消后继续等待。

参考答案:完整取消实验
stream-cancel.mjs
import { Writable } from 'node:stream';
import { pipeline } from 'node:stream/promises';
import { setTimeout as delay } from 'node:timers/promises';
import assert from 'node:assert/strict';

const controller = new AbortController();
let released = false;
let count = 0;
async function* tasks({ signal }) {
  try {
    for (let id = 1; id <= 10; id += 1) {
      await delay(5, undefined, { signal });
      yield { id };
    }
  } finally { released = true; }
}
const sink = new Writable({
  objectMode: true,
  write(task, encoding, callback) {
    count += 1;
    if (count === 3) controller.abort();
    callback();
  }
});
await assert.rejects(pipeline(tasks, sink, { signal: controller.signal }), { name: 'AbortError' });
assert.equal(count, 3);
assert.equal(released, true);
console.log('第三项取消,生成器资源已释放');

执行 node stream-cancel.mjs,预期通过。这里对象模式只在内存中传递小任务对象,并不意味着可以把整个文档对象无限塞入流。生成器 finally 在本例只修改观察标记,真实实现可以关闭自己拥有的资源,但仍要处理清理失败与重复清理的边界。

读取模式与异步迭代如何选择#

Readable 可以通过 data 事件进入持续流动的消费方式,也可以通过显式读取或异步迭代逐步消费。它们不是可以任意混合的几套独立订阅:同时使用多种消费方式,会让数据究竟被谁拿走变得难以推理。为一条流选择一种主要消费协议,组合时再通过明确适配器连接。

for await 很适合逐块执行业务处理,因为循环体中的 await 自然表达“处理完本块再继续取下一块”。但如果循环体只启动一个异步保存函数而不等待,仍会产生无限并发,失去顺序和容量约束。异步迭代提供的是取得下一项的协议,不会猜测你在循环体里发起的其他任务。

提前退出循环也涉及资源生命周期。默认流异步迭代在某些提前结束情况下会销毁流,因此一个只想先看前几个字节的函数,不应假定之后另一个调用方还能继续消费同一对象。需要窥探或复用时,应明确选择对应的迭代选项、缓冲策略或重新打开资源,并核对所用 Node 版本的契约。

转换器如何处理跨块业务记录#

假设文档采用一行一条 JSON 的格式。转换器除了 UTF-8 解码器,还需要一个尚未结束的行缓冲。每来一块,把完整文本接到尾部,提取所有完整行,留下最后的不完整行。输入结束时还要决定无换行的最后一行是否有效。若忘记最终冲刷,最后一条记录可能悄悄丢失。

这个尾部缓冲也必须有限。攻击者可以一直发送不带换行的内容,让所谓“逐行流式解析”不断积累一行。总字节上限和单行上限承担不同职责:前者限制整份输入,后者限制解析器一次需要保留的记录。真正有界的设计,需要把每一块持久状态都纳入容量分析。

转换还可能放大数据。压缩文件很小,解压后却很大;一条输入记录也可能展开成许多输出记录。因此只限制上传原始字节不一定够,还要在解压或展开之后限制输出规模、记录数量和处理时间。背压能控制速度差,不能自动阻止无限放大或业务上不合理的总量。

文件发布需要与字节传输分开设计#

团队助手保存文档时,可以先为本次上传创建随机临时文件,把完整管道写入它,完成校验后再将它登记为可见文档。这样列表接口不会看到只写了一半的内容。失败时只清理本次临时文件,避免把其他任务已经发布的结果误删。临时文件路径必须由服务端生成,不能直接把用户文件名当成任意磁盘路径。

即使写入管道完成,也不应扩大成“发生断电时保证永久保存”的承诺。操作系统缓存、同步刷新和目录元数据持久性是更深一层的存储问题。应用需要什么级别的耐久性,应由数据库或存储方案的明确契约承担。基础实验验证的是数据经过流处理完成,不是模拟断电恢复。

下载输出则相反:响应头一旦发出,后面发生的读取失败通常无法再替换为一个普通 JSON 错误。可以在开始响应前完成必要权限与资源检查,但仍必须接受中途网络断开。前端要能区分完整文件与截断结果,必要时使用长度、摘要或可恢复下载协议。这说明流处理与 HTTP 协议不能只靠一行 pipe 就被完全解释。

观察内存与调试停滞#

管道停住时,先检查自定义 write 或 transform 是否在每个分支都调用了 callback,再检查是否等待一个永远不会到来的事件。错误地在 write 返回 false 后重写同一块,会导致重复数据;忽略 false 则可能让缓冲持续增加。两种错误都可能在小文件里看不出来,所以实验刻意把阈值设得很小。

可以观察 writableLength、readableLength 以及相应 highWaterMark,了解当前内部缓冲,而不是只看进程总内存。总内存还包含运行时、其他请求和外部缓冲,短时间不上升也不能证明永远有界。用不同输入规模比较峰值趋势,比只执行一次十字节样本更接近实际容量判断。

旧资料中的手工 pipe 与事件清理仍能帮助理解机制,现代代码则可以使用 Promise 版 pipeline、finished 和 AbortSignal 把共同生命周期表达得更集中。变化的是组合工具,仍然成立的原则是生产者必须接受容量反馈、每个资源有关闭路径、每个失败能到达拥有该任务的边界。

在服务里验证流时,还应同时发起一个很小的健康请求,观察大文件处理是否影响它的响应。总吞吐高但小请求长时间排队,说明单次转换或并发策略仍可能不合理。把吞吐、尾部延迟、峰值内存和取消后的资源数量一起观察,才能判断这条管道是否适合真实交互负载。

另一个验收点是重复运行:同一任务工厂连续创建多条管道后,监听器、句柄和临时文件应回到预期基线。单次成功无法暴露长期积累问题,而长期服务恰好会反复经过这些路径。

每次都应验证关闭。

本章实际验证范围#

三个文件实际通过:慢写入累计二十字节并暂停五次,超限管道拒绝且目标销毁,第三项取消后生成器 finally 执行。未进行大文件压力或实际网络上传测试。

验收、自测与参考#

验收包括慢写入总字节不重复,超限管道拒绝且相关流销毁,取消时生成器 finally 执行。自测一:write 返回 false 是否应重写当前块?答案:不应,当前块已被接受,应等待后续容量。自测二:highWaterMark 是否限制整个进程内存?答案:不是,只是某个缓冲的反馈阈值。自测三:pipeline 失败是否会自动撤销已经写出的文件内容?答案:不会,发布与回滚需要业务协议。

完整契约见 Stream 官方文档,实际背压影响见 Backpressuring in Streams,可取消等待见 Timers Promises

原有课程整理于 2026-09-10;Node / Electron 扩充于 2026-09-11。示例环境与验证范围以正文为准。
原创中文学习手册,阅读结构参考 Vue 文档;非 Vue 官方教材。
下载本章 Markdown

支持中文和英文全文搜索 · ↑ ↓ 选择 · Enter 打开 · Esc 关闭