代表性面试主题

后端面试:如何用 Node.js stream.compose 设计可取消的数据流水线?

后端困难
Offer.cc 编辑团队发布 更新

题干

一个 Node.js 服务要把上传、解压、校验和写入对象存储串成流水线。如何用 stream.compose 保证背压、错误传播、取消和资源清理?

题干与适用场景

一个 Node.js 服务接收大文件上传,需要依次完成解压、病毒扫描、格式校验和对象存储写入。旧实现手动监听 data 事件,偶发内存上涨、客户端断开后任务仍继续,以及中间步骤失败时临时文件没有清理。

请使用 stream.compose() 或等价的 pipeline 组合方式,说明可读、可写和 Transform 阶段如何连接,如何利用背压、AbortSignal、错误传播和最终清理保证一次任务可控。

面试官考察点

  • 能否区分 pipepipelinecompose 的边界与生命周期。
  • 能否解释背压如何限制生产速度,而不是简单提高队列上限。
  • 能否让取消、异常和客户端断开沿整条流水线传播。
  • 能否处理异步生成器、资源释放、幂等写入与可观测性。

回答前需要澄清的问题

  1. 输入来自 HTTP request、文件还是对象存储 SDK?输出是否支持分片和重试?
  2. 每个阶段是 Node stream、Web Streams、AsyncIterable 还是普通函数?
  3. 病毒扫描和校验是否会产生外部进程、临时文件或数据库状态?
  4. 取消后对象存储是否允许中止 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 语义处理相邻阶段。示例:

js
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 是否有界,并确认吞吐接近最慢阶段而非无限缓存。

取消后能否直接重试?

先确认所有资源已关闭、临时状态可识别,再从幂等且可恢复的边界重试。不可中断的外部写入应先完成补偿或标记未知状态。

公开来源

同类题目