后端面试:如何用 Node.js stream.compose 设计可取消的数据流水线?
题干与适用场景
一个 Node.js 服务接收大文件上传,需要依次完成解压、病毒扫描、格式校验和对象存储写入。旧实现手动监听 data 事件,偶发内存上涨、客户端断开后任务仍继续,以及中间步骤失败时临时文件没有清理。
请使用 stream.compose() 或等价的 pipeline 组合方式,说明可读、可写和 Transform 阶段如何连接,如何利用背压、AbortSignal、错误传播和最终清理保证一次任务可控。
面试官考察点
- 能否区分
pipe、pipeline和compose的边界与生命周期。 - 能否解释背压如何限制生产速度,而不是简单提高队列上限。
- 能否让取消、异常和客户端断开沿整条流水线传播。
- 能否处理异步生成器、资源释放、幂等写入与可观测性。
回答前需要澄清的问题
- 输入来自 HTTP request、文件还是对象存储 SDK?输出是否支持分片和重试?
- 每个阶段是 Node stream、Web Streams、AsyncIterable 还是普通函数?
- 病毒扫描和校验是否会产生外部进程、临时文件或数据库状态?
- 取消后对象存储是否允许中止 multipart upload,任务是否需要恢复?
30 秒回答
我会把每个阶段定义成可读、可写或 Transform/AsyncIterable,并用 stream.compose() 组合成一个 Duplex,再由 pipeline 驱动到最终目的地。生产者只在下游可接受时继续生成,让背压限制内存。把请求断开、超时和业务取消合并为一个 AbortSignal,传给支持 signal 的阶段;任何错误都让整条链失败。临时文件、子进程和 multipart upload 在 finally 或 abort handler 中清理,写入使用幂等键。指标记录吞吐、队列长度、峰值 RSS、取消原因和清理结果。
分步骤深入解答
先定义每个阶段的契约
明确输入输出类型、chunk 大小、是否允许 null、是否会阻塞,以及阶段完成时谁拥有资源。异步生成器必须消费 source 并按需 yield;普通函数不能偷偷把整文件读进内存。
用 compose 组合可复用阶段
compose 把多个 stream、AsyncIterable 或函数连接成新的 Duplex;它内部按 pipeline 语义处理相邻阶段。示例:
import { compose } from 'node:stream';
async function* validate(source) {
for await (const chunk of source) {
checkChunk(chunk);
yield chunk;
}
}
const processing = compose(decompress(), validate, scan());实际代码还要把 processing 接到目标 writable,并统一监听完成与错误,而不是每段各自吞掉异常。
让背压成为默认控制面
下游 writable 未准备好时,Readable/Transform 应暂停继续生产。不要在 data 回调里无界 push,也不要只靠增大 high water mark 掩盖消费者变慢。压测时记录每阶段队列、吞吐和 RSS,确认最慢阶段决定整体速度。
传播取消与客户端断开
把请求断开、deadline、人工取消合并到一个 AbortController。将 signal 传给支持它的 compose 阶段和外部 SDK;abort 后停止读取、销毁下游并等待关闭完成。取消是正常结束原因,日志需区分 abort、业务失败和网络错误。
统一错误与资源清理
让一个拥有者调用 pipeline/compose,统一接收第一个错误并触发销毁。临时文件、扫描进程、socket 和 multipart upload 必须在成功、失败和取消三条路径都释放;清理失败要告警并带任务 ID,不能覆盖原始错误。
设计幂等与观测指标
用上传 ID 和阶段版本生成幂等键;对象存储写入完成后再提交业务状态。记录输入字节、输出字节、处理时长、峰值 RSS、取消原因、错误阶段和清理耗时。重试只从可恢复边界重新开始,避免重复写入。
高质量示范回答
我会先把上传、解压、扫描、校验和写入定义为有明确输入输出与资源所有权的阶段,再用 compose 组合可读、可写或异步迭代器,最终交给统一的 pipeline 驱动。背压由下游消费速度控制,禁止在 data 事件中无界缓存。请求断开、超时和人工取消共享一个 AbortSignal,并传递给支持 signal 的阶段和存储 SDK。统一的拥有者处理首个错误,销毁整条链;临时文件、扫描进程和 multipart upload 在成功、失败、取消路径都清理。写入用幂等键,指标覆盖吞吐、队列、RSS、取消原因和清理结果,重试只从安全边界恢复。
常见错误
- 把整条流收集成 Buffer,再声称使用了 compose。
- 只监听最终 writable 的
error,遗漏异步生成器或外部进程错误。 - 让客户端断开后源流继续读取并写入对象存储。
- 用无限增大 high water mark 解决背压问题。
- 只在成功路径删除临时文件,忽略 abort 和异常。
- 重试没有幂等键,导致对象重复或业务状态重复提交。
追问及应对
compose 和 pipeline 如何分工?
compose 负责把阶段组合成可复用的 Duplex;pipeline 负责驱动端到端连接、传播错误并等待关闭。两者可以组合使用,但所有权和错误处理应集中。
异步生成器抛错会发生什么?
错误应让组合流失败并销毁相邻阶段。调用方仍需等待 pipeline 完成清理,并记录实际失败阶段,不能只依赖进程未退出。
如何验证背压真的生效?
用比消费者更快的输入和受控慢写入压测,观察队列长度与 RSS 是否有界,并确认吞吐接近最慢阶段而非无限缓存。
取消后能否直接重试?
先确认所有资源已关闭、临时状态可识别,再从幂等且可恢复的边界重试。不可中断的外部写入应先完成补偿或标记未知状态。