在 Buffer 与字符编码中,我们讨论了字节怎样被保存、切分和还原。面对大文件或持续到来的数据,还有一个问题:即使每次只处理一小块,如果下一块来得太快,尚未处理完的数据应该放在哪里?
假设生产者已经提交了八条记录,消费者却刚开始处理第一条。两边都使用异步 API,JavaScript 也没有在原地等待,但另外七条记录仍然需要留在某个队列里。如果生产一直快于消费,这个队列就会持续增长。
Stream 将连续的数据组织成可以逐步生产、交付和消费的序列;背压则把下游的处理能力反馈给上游,让暂时处理不过来的数据不再被无限提交。 本文沿着编号为 1~8 的数据,观察它们在哪里等待、由谁推动,以及暂停以后怎样继续前进。
本文以 Node.js v24.16.0 为验证环境,使用原生 ESM。带有 .mjs 文件名的代码块都是完整、独立的实验,保存后使用 node 文件名.mjs 运行。实验不连接数据库或外部服务;文件实验只创建并清理自己的临时目录。定时器用于制造异步等待,具体耗时不作为结论。
1. 异步运行,为什么仍然会积压
1.1 操作返回以后,工作可能还没有完成
异步 API 让调用方可以在操作尚未完成时继续执行。它解决了等待期间怎样利用执行机会的问题,却没有保证新提交的工作恰好能被及时完成。
把同一种数据单位放在一条不丢弃、不复制记录的处理链里,可以先建立一个计数关系:
累计生产数量 - 累计完成数量 = 当前在途数量
在途数据包括正在处理的部分,也包括已经生产但尚未开始处理的部分。如果每秒生产 100 条、完成 10 条,在没有减速或其他处理策略的情况下,每秒就会多出约 90 条未完成记录。
分块处理使应用能够在完整输入到达之前开始工作,也避免必须一次保留全部输入。但要获得这个收益,生产速度还需要受到消费速度的约束。如果一边分块读取,一边把所有块放进数组等以后再处理,内存仍然随输入增长。
图中的分块处理允许少量预读和等待。缓冲能够吸收短期速度波动;持续供大于求时,则需要进一步限制进料速度。
1.2 先观察一次没有遵守反馈的写入
下面的 Writable 是一个可写流:调用方交给它一条记录,它经过一次异步等待后报告处理完成。objectMode 表示按 JavaScript 值传递数据,此时阈值按项数计量。
// unchecked-writes.mjs
import assert from "node:assert/strict";
import { Writable } from "node:stream";
import { finished } from "node:stream/promises";
let completed = 0;
const received = [];
const sink = new Writable({
objectMode: true,
highWaterMark: 2,
write(value, encoding, callback) {
setTimeout(() => {
received.push(value);
completed++;
callback();
}, 5);
},
});
for (let value = 1; value <= 8; value++) {
const canContinue = sink.write(value);
console.log(value, canContinue, sink.writableLength);
}
console.log("提交结束,已完成", completed);
sink.end();
await finished(sink, { cleanup: true });
assert.deepEqual(received, [1, 2, 3, 4, 5, 6, 7, 8]);
console.log("处理结束,已完成", completed);
输出:
1 true 1
2 false 2
3 false 3
4 false 4
5 false 5
6 false 6
7 false 7
8 false 8
提交结束,已完成 0
处理结束,已完成 8
同步循环提交八条记录期间,定时器回调尚未获得执行机会。因此,提交结束时,八条记录全部处于未完成状态。writableLength 统计待完成写入的长度,包含正在处理但尚未报告完成的那一项。
这个实验的输入有限、处理始终成功,用于观察排队,并不是内存压力测试。返回值从第二项开始变成 false,后续数据却仍然被接收。理解这个反馈,要先看清一块数据经过的几个位置。
2. 一块数据怎样穿过 Stream
2.1 数据源、流对象与底层目标
Readable 提供可读数据,Writable 接收写入的数据。它们分别把底层资源与使用者连接起来:文件、socket 或按需生成器负责提供数据,文件写入、网络发送或业务函数负责处理数据。
这些位置保存的数据处于不同阶段:
- Readable 缓冲中的数据已经取得,但尚未交给消费者。
- Writable 中的数据已经提交给写入端,但尚未完成它定义的处理。
- 底层系统还可能有自己的缓冲或在途操作。
一个 chunk 是一次交付的数据块。在常见字节流中,它通常是 Buffer;开启文本解码后可以是字符串,对象模式下也可以是其他 JavaScript 值。队列可能只是保存 chunk 的引用,并不要求每次交付都复制全部字节;转换或编码则可能另外分配存储。
chunk 的边界由数据产生与交付方式决定,不天然对应一行文本或一条消息。中文字符可能跨块,完整字符也可能只组成半条 JSON;字符解码和消息分帧的责任,已在 Buffer 篇中分别讨论。
2.2 使用者提交需求,实现者连接资源
Stream 把使用者接口和实现者接口分开。消费 Readable 时,可以使用异步迭代、pipe(),或在合适的读取模式下调用 read();实现者提供 _read(),在流需要补充数据时从底层取数,再用 push(chunk) 交给流。
写入端也有两层。使用者调用 write(chunk);Writable 管理排队,并调用实现者的 _write(chunk, encoding, callback)。实现者在当前块处理完成后调用 callback(),失败时调用 callback(error)。构造选项里的 read()、write() 分别是提供这些实现方法的简写,应用不应直接调用带下划线的方法。
一次逻辑上的推进可以写成:
产生消费需求 → _read() 补充数据 → push(chunk)
→ 交付 chunk → write(chunk) → _write() 处理
→ callback() 报告当前处理完成
实际执行会包含预读、同步交付和内部调度。某些条件下,新来的 chunk 可以直接交给等待它的消费者,不必先在可读队列中停留。也不能把这条链画成“每经过一个箭头,就等待一轮事件循环”。
2.3 事件循环推动完成通知,流协议协调速度
异步文件操作通常通过 libuv Worker Pool 执行,网络 I/O 通常依赖非阻塞 socket 与平台通知机制;结果就绪或操作完成后,相应回调才有机会继续推动 JavaScript。具体路径可参照事件循环与异步 I/O。
Stream 对这些资源提供统一的数据交付和流量协议。创建一个流不会为它分配专属 JavaScript 线程。自定义转换中的同步计算仍然占用当前线程;即使队列受控,一次计算过久也会拖延其他回调。
对外消费时,Readable 可以持续通过 data 事件交付,也可以由消费者主动取得数据。同一条流宜选择一种主要消费方式,避免混用 data、readable、pipe() 和异步迭代,使数据归属及推进时机难以判断。
3. 下游怎样让上游停下来
3.1 返回 false 后,为什么第三项仍然能写进去
回到第一组实验。阈值 highWaterMark: 2 表示待完成数量达到相应条件时,开始要求调用方暂停提交。当前写入已经被接收,返回值影响的是后续是否应当继续写入。
对于这个正常可写、异步处理的目标,前几次调用对应如下状态:
| 调用 | 正在处理 | 等待处理 | 待完成数量 | 返回值 |
|---|---|---|---|---|
write(1) |
1 | 空 | 1 | true |
write(2) |
1 | 2 | 2 | false |
write(3) |
1 | 2、3 | 3 | false |
false 表示当前数据已经接收,后续提交应暂停。 重写当前 chunk 会导致重复;忽略返回值继续写,则可能持续增长队列。它也不是一次写入成功完成的证明,真正的错误需要通过相应错误处理路径观察。
因此,highWaterMark 是反馈阈值,而不是硬性拒收上限。一个很大的 chunk 就可能越过阈值;流已经出错或销毁时,又有另外的失败语义,不能只凭布尔值推断原因。此版本的写入契约明确区分了继续提交与实际完成。
3.2 暂停交付之后,可读侧也要停止补充
手动向 Writable 写入时,调用方负责响应返回值。通过 pipe() 连接时,连接逻辑会在目标 write() 返回 false 后暂停源的流动,并在目标恢复后继续。
下面只是单源、单目标连接的机制示意,省略了结束、错误、销毁和监听器清理,不能替代完整的 pipeline():
source.on("data", (chunk) => {
if (!destination.write(chunk)) source.pause();
});
destination.on("drain", () => source.resume());
pause() 先暂停的是 Readable 持续通过 data 向外交付的过程。已经预读的数据仍留在缓冲中,已经发出的底层操作也不保证立即取消。
如果消费者不再取走数据,可读缓冲就难以继续腾出空间。Readable 实现者调用 push(chunk) 得到 false 后,应停止继续推送,等待后续 _read() 请求再恢复;外部主动推送的数据源也需要在这里接入自己的暂停机制。Readable 的这一侧没有等待自身 drain 的协议。
因此,暂停通过两道反馈逐步传递:
目标待完成数据增多 → write() 返回 false → 暂停向目标交付
→ 可读缓冲暂时消耗不掉 → push() 返回 false → 暂缓继续生产
Node.js 会根据需求与缓冲状态决定何时再请求数据,但任意自定义代码若在 push() 返回 false 后仍持续推送,依然可以让缓冲增长。背压需要每一段都保留并响应反馈。
3.3 阈值的单位决定了它能约束什么
普通字节流通常按字节计量;对象模式按项数计量。两个小整数与两个引用大 Buffer 的对象,在对象模式里都算两项,但占用的内存并不相同。
在非对象模式下,调用 readable.setEncoding() 后,可读缓冲按字符串长度计量,也就是 UTF-16 码元数,不能直接当成 Buffer 字节数。阈值不保证每个 chunk 的大小,更不等于整个进程或操作系统的内存上限。默认值还可能随版本、平台或具体流实现变化,本文实验因此都显式设置关键阈值。
4. 下游处理完成以后,数据怎样恢复
4.1 当前一项完成,与可以继续提交
对普通 _write() 实现,Node.js 会等待当前写入通过回调报告完成,再安排后续处理。连续调用多次 write() 主要增加待完成量,不会自动把底层处理变成同等数量的并行任务;支持批量写入的 _writev() 是另一种实现能力。
callback() 说明当前块处理完成,drain 则说明此前要求等待的写入端恢复了继续提交条件。不是每完成一项都发出 drain。
在本文版本的正常持续写入路径中,内部已经记录“需要 drain”,且待完成长度降到 0 时,才会发出该事件。因此,不能把它改述为“一低于 highWaterMark 就触发”。流正在结束、出错或销毁时,也不应期待一定会再收到 drain;任务完成应通过 pipeline() 或 finished() 等接口观察。对应版本的 Writable 实现可以追踪这段状态更新。
等待期间,JavaScript 执行权会交还给运行时,使处理完成的通知有机会运行。用同步循环反复检查 writableLength,反而可能占住执行完成回调所需的线程。
4.2 让同一组数据经历多轮暂停与恢复
下面把按需生产的 Readable 与慢速 Writable 连接起来。生产者只在 read() 请求中生成数据,看到 push() 返回 false 就返回。消费者每次完成一项后才调用回调。
// flow-and-pressure.mjs
import assert from "node:assert/strict";
import { Readable, Writable } from "node:stream";
import { pipeline } from "node:stream/promises";
let produced = 0;
let completed = 0;
const received = [];
const source = new Readable({
objectMode: true,
highWaterMark: 2,
read() {
while (produced < 8) {
if (!this.push(++produced)) return;
}
this.push(null);
},
});
const sink = new Writable({
objectMode: true,
highWaterMark: 2,
write(value, encoding, callback) {
setTimeout(() => {
received.push(value);
completed++;
callback();
}, 5);
},
});
function snapshot(label) {
console.log(
label,
produced,
completed,
source.readableLength,
sink.writableLength,
);
}
console.log("状态 已生产 已完成 可读缓冲 待完成写入");
source.on("pause", () => snapshot("暂停"));
sink.on("drain", () => snapshot("drain"));
await pipeline(source, sink);
assert.deepEqual(received, [1, 2, 3, 4, 5, 6, 7, 8]);
snapshot("完成");
在本文版本中输出:
状态 已生产 已完成 可读缓冲 待完成写入
暂停 4 0 2 2
drain 4 2 2 0
暂停 6 2 2 2
drain 6 4 2 0
暂停 8 4 2 2
drain 8 6 2 0
暂停 8 6 0 2
完成 8 8 0 0
第一次暂停时,可以把数据位置展开为:
| 所在位置 | 数据 |
|---|---|
| 尚未生产 | 5、6、7、8 |
| Readable 缓冲 | 3、4 |
| Writable 正在处理 | 1 |
| Writable 等待处理 | 2 |
| 已完成 | 空 |
当 1、2 完成后,目标发出 drain。3、4 可以继续交付,Readable 再根据需求补充 5、6。上游暂停期间,下游始终可以消化当前数据;恢复时,也能利用已有预读,而不必每次都重新等待底层取得第一块数据。
这个实验中的“已生产减去已完成”没有超过 4,但这只是当前实现、输入和观察过程的结果。预读与取出可能交错,某个函数内部的瞬时长度还可能越过阈值,不能把日志推广为所有流通用的精确内存上限。
最后一次暂停之后没有再打印 drain:源已经交付完全部数据,写入端进入结束流程。pipeline() 等待整个任务完成,避免把“恢复继续写”误当成“所有工作结束”。
5. 压力怎样穿过转换与业务处理
5.1 Transform 把输出消费能力传回输入
双工流 Duplex 同时具备可读和可写能力,例如 socket 的接收与发送。Transform 是一种输入经过处理后产生输出的双工流,例如压缩器。它同时维护可写侧和可读侧,输入与输出的大小也可能不同。
考虑文件压缩后发送给慢客户端:
文件 Readable → gzip Transform → 网络 Writable
当网络端要求暂停时,压缩结果暂时无法继续交付。如果 Transform 仍无限处理输入,积压就会转移到输出缓冲。因此,输出侧的压力必须限制输入侧继续推进。
实现者在 _transform(chunk, encoding, callback) 中处理输入,通过 push() 或回调的结果参数提供输出,并在当前输入处理完毕后调用回调。在本文版本的一般输出积压路径中,Transform 内部会暂缓完成底层写入回调,等可读侧出现新的读取需求再推进。这样,输出消费能力便能继续影响上游输入。Transform 实现展示了这两个方向的协调。
这种协调不能自动打断任意 JavaScript 循环。如果一次转换内部展开出大量结果,并在 push() 返回 false 后继续推送,单次转换仍可能制造很大的积压。遇到一条输入产生海量输出的任务,还需要设计可暂停的分批产出过程,约束单次工作量与中间状态。
5.2 异步监听器可能把积压搬到另一个队列
另一种更隐蔽的断点发生在业务处理处:
// 错误假设示意:这里的 await 不会使 data 分发等待保存完成。
source.on("data", async (record) => {
await saveRecord(record);
});
EventEmitter 调用监听器,但不会等待监听器返回的 Promise 再发出下一次事件。于是,可读缓冲可能很短,保存请求却已经大量并发启动。队列没有消失,只是搬到了业务函数、数据库驱动或其他任务系统里。
下面用有限输入比较事件消费与逐项等待。事件版本把任务 Promise 留下来,是为了在实验结束前等待并观察它们;它没有在交付期间限制任务数量。
// async-consumers.mjs
import assert from "node:assert/strict";
import { Readable } from "node:stream";
import { finished } from "node:stream/promises";
import { setImmediate as nextTurn } from "node:timers/promises";
function* records() {
for (let value = 1; value <= 8; value++) yield value;
}
async function run(mode) {
let active = 0;
let peak = 0;
async function processRecord(value) {
active++;
peak = Math.max(peak, active);
await nextTurn();
active--;
return value;
}
const source = Readable.from(records(), { highWaterMark: 2 });
if (mode === "event") {
const pending = [];
source.on("data", (value) => pending.push(processRecord(value)));
await finished(source, { cleanup: true });
await Promise.all(pending);
} else {
for await (const value of source) {
await processRecord(value);
}
}
return peak;
}
const eventPeak = await run("event");
const iteratorPeak = await run("iterator");
assert.equal(eventPeak, 8);
assert.equal(iteratorPeak, 1);
console.log("事件消费的未完成任务峰值", eventPeak);
console.log("逐项等待的未完成任务峰值", iteratorPeak);
输出分别为 8 和 1。这里的异步迭代消费的是具备流量协议的 Readable,并且循环体真正等待了处理完成。普通 EventEmitter 包装成异步迭代器,不保证生产者会跟随减速;这一点可对照 EventEmitter 篇。
5.3 让业务的真实完成进入写入协议
如果业务更适合接入处理链,可以把保存操作包装为 Writable。下面是接入示意,其中 saveRecord 应返回代表当前记录实际处理完成的 Promise:
const sink = new Writable({
objectMode: true,
highWaterMark: 2,
write(record, encoding, callback) {
Promise.resolve()
.then(() => saveRecord(record))
.then(() => callback(), (error) => callback(error));
},
});
保存完成以后才调用回调,未完成业务就会反映为未完成写入。反之,如果只是把记录扔进另一个无限数组,然后立即调用回调,Stream 就会误以为这部分工作已经完成。
阈值为 2 也不表示并行执行两次保存。这个普通 Writable 按顺序处理;确实需要并发时,要另外限制运行中任务数与等待队列,并定义结果顺序、错误和取消行为。
逐项异步迭代中也有一个相似陷阱:经典 Node.js Stream 的 write() 返回布尔值,因此 await destination.write(chunk) 不会等待写入完成。处理两端传输时,优先把协调交给完整的管道接口。
6. 一条处理链怎样可靠地结束
6.1 pipeline 连接流量、完成与错误
pipe() 能协调相邻流之间的背压;pipeline() 进一步集中处理整条链路的完成、错误传播及流资源销毁。使用 Promise 版本,就可以在任务入口等待最终结果。
下面创建一份有限文本,逐块读取、压缩并写入临时文件,再验证恢复出的内容。校验阶段有意读回这份小样本;实际大文件可以改为逐块摘要比较,不必把全部内容读回内存。
// compress-file.mjs
import assert from "node:assert/strict";
import { createReadStream, createWriteStream } from "node:fs";
import { mkdtemp, writeFile, readFile, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { pipeline } from "node:stream/promises";
import { promisify } from "node:util";
import { createGzip, gunzip } from "node:zlib";
const directory = await mkdtemp(join(tmpdir(), "stream-backpressure-"));
try {
const input = join(directory, "input.txt");
const output = join(directory, "input.txt.gz");
const original = Buffer.from("Guang,你好👋\n".repeat(4096), "utf8");
await writeFile(input, original, { flag: "wx" });
await pipeline(
createReadStream(input, { highWaterMark: 4096 }),
createGzip(),
createWriteStream(output, { flags: "wx" }),
);
const restored = await promisify(gunzip)(await readFile(output));
assert.deepEqual(restored, original);
console.log("逐块压缩完成,内容校验通过");
} finally {
await rm(directory, { recursive: true, force: true });
}
流式处理没有要求等整个文件读取完才开始压缩。下游变慢时,压力可以逐段向前传递;出错时,任务入口也能得到失败结果。
6.2 完成、关闭与业务确认属于不同层次
有几个名称相近的信号,需要放回各自的位置理解:
| 操作或信号 | 含义 |
|---|---|
Readable 的 push(null) |
实现者声明不会再产生数据,已经缓冲的数据仍需消费 |
Readable 的 end |
可读数据已经消费完 |
writable.end() |
调用方声明不再提交后续写入 |
Writable 的 finish |
结束写入后,待完成数据及最终处理已按该流的协议完成 |
close |
流及其底层资源关闭;单独收到它不能判断任务成功 |
这些信号都不能自动替代远端业务确认或文件断电持久性的保证。写入端报告完成,可能只是数据已经交给相应底层系统;如果业务需要确认对方已经保存,仍需要应用协议。文件输出中途失败,也可能留下部分内容;需要结果原子可见时,可以设计临时文件与完成后的替换流程。
6.3 错误与取消要沿链路传播
下面在遇到第 3 项时,分别注入处理错误和发出取消信号。每次使用新的源和目标,不复用失败后的流。
// pipeline-stop.mjs
import assert from "node:assert/strict";
import { Readable, Writable } from "node:stream";
import { pipeline } from "node:stream/promises";
async function run(mode) {
const controller = new AbortController();
const completed = [];
const source = Readable.from([1, 2, 3, 4, 5, 6, 7, 8], {
highWaterMark: 2,
});
const sink = new Writable({
objectMode: true,
highWaterMark: 2,
write(value, encoding, callback) {
if (value === 3) {
if (mode === "error") callback(new Error("模拟第 3 项失败"));
else {
controller.abort();
callback();
}
return;
}
setImmediate(() => {
completed.push(value);
callback();
});
},
});
await assert.rejects(
pipeline(source, sink, { signal: controller.signal }),
mode === "error"
? { message: "模拟第 3 项失败" }
: { name: "AbortError" },
);
assert.deepEqual(completed, [1, 2]);
assert.equal(source.destroyed, true);
assert.equal(sink.destroyed, true);
console.log(mode, "已完成", completed.join(","), "两端已销毁");
}
await run("error");
await run("abort");
输出:
error 已完成 1,2 两端已销毁
abort 已完成 1,2 两端已销毁
这里的取消时机经过刻意控制,取消时没有额外未完成的模拟业务操作。真实数据库请求、网络调用或自定义定时器是否停止,需要这些操作响应信号,并由实现者在销毁路径清理资源。销毁流不会自动回滚已完成的 1、2,也不会让任意外部 Promise 自动取消。
pipeline() 适合调用方拥有并管理的处理链。对于共享资源或 HTTP 响应,要预先考虑销毁会影响谁;链路错误可能关闭 socket,不能假定失败后仍能补发一份完整的错误响应。
7. 浏览器怎样表达相同的反馈
浏览器也面对输入快于处理的问题。Web Streams 使用可读流、可写流和转换流组织数据,主要反馈形式如下:
| 关注点 | 经典 Node.js Stream | Web Streams |
|---|---|---|
| 提交写入 | write(chunk) 返回能否继续提交的布尔值 |
writer.write(chunk) 返回表示这次写入完成或失败的 Promise |
| 恢复提交 | 遵守 false,等待 drain |
观察 desiredSize,通过 writer.ready 等待恢复条件 |
| 补充可读数据 | _read() 与 push() 返回值 |
pull(controller)、enqueue() 与 controller.desiredSize |
| 连接处理链 | pipeline() |
pipeThrough()、pipeTo() |
对正常可写的 Web 流,desiredSize 可以理解为阈值减去队列总尺寸。它可以为负;即使返回 Promise,调用方连续不等待地写入,仍然可以让队列超过阈值。writer.ready 在正常背压恢复时关注余量重新为正,不要求当前队列完全为空,也不表示某次写入已经完成。因此,它与经典 Node.js 的 drain 不能机械互换。Streams 标准定义了这些反馈关系。
enqueue() 也不返回 Node.js push() 那样的布尔值。按需数据源可以把生产放进 pull();外部推送源则需要结合 desiredSize 接入实际暂停能力。普通 Web 流默认按 chunk 数量计量,可以用排队策略改为按字节计量;专用可读字节流另有对应规则。
两套接口的差别不等于浏览器与 Node.js 的能力分界,现代 Node.js 同样支持 Web Streams。无论选择哪套接口,都要让业务的真实完成进入消费协议;传统 WebSocket、普通事件或 Worker 消息,也不会仅因为处理函数带有 async 就自动获得端到端流量控制。
8. 怎样判断背压真正生效
背压不要求上游时刻不停地生产。快的一端补充一批数据后等待,慢的一端消化后再继续,这种间歇是协调速度的一部分。如果在途数据能够保持受控,长期平均进料速度也就必须适应完成速度。
实际排查时,可以从现象沿着数据位置追问:
| 现象 | 优先检查 |
|---|---|
writableLength 持续增长 |
是否忽略 write() 返回的 false;写入回调是否迟迟未完成 |
readableLength 持续增长 |
生产者是否忽略 push() 返回的 false;底层推送是否能够暂停 |
| 流缓冲不大,业务任务却越来越多 | 是否在事件监听器里无限启动 Promise,或把任务转存进其他队列 |
| 队列不增长,但其他请求明显变慢 | 单块转换是否包含过重的同步计算 |
| 上游不再读取,内存仍然很高 | 是否提前加载全部输入、持有已处理数据,或存在其他缓冲层 |
| 取消以后仍有副作用 | 外部操作是否支持取消,销毁路径是否释放自己的资源 |
观察时,要明确计量单位和所有权:Writable 的长度已经包含正在处理的那部分,不能重复相加;对象数量也不能直接换算成字节内存。heapUsed、Buffer 底层存储和进程 rss 的范围不同,单个计数稳定并不足以证明整个进程内存稳定。
对能够暂停的数据源,背压把处理能力逐段传回生产端。对无法减速的来源,应用必须另外选择有限缓存、丢弃、采样、落盘或拒绝输入等策略。长期输入速度超过处理速度时,有限存储、无限接收和永不丢失无法同时满足。
回看编号为 1~8 的记录:数据先被生产,再被交付,最后完成处理;这三个时刻之间的差额,就是应用必须负责的在途工作。沿着队列找到等待的数据,再沿着返回值、完成回调和恢复通知找到反馈,才能判断一条处理链是否真正受到下游能力的约束。