You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在createReadStream的data事件中正确使用fs.write拼接二进制文件

问题根因
  • 直接在data事件回调中调用write()未等待写入完成,也没有处理流背压:当写入流内部缓冲区满时,write()会返回false,此时继续写入会导致数据堆积、丢失,最终写入字节数异常、运行结果不固定。
  • 即使将data事件回调改为async函数使用await,也会导致多个写入操作并发执行,无法保证数据顺序,反而会加剧问题。
修复方案

方案1:使用异步迭代器(推荐,符合现代JS语法习惯)

直接用for await...of迭代读取流,自动异步获取chunk,同时手动处理背压即可:

const { createReadStream, createWriteStream } = require('fs')

async function concatFile (filename, writeStream) {
  const readStream = createReadStream(filename, { highWaterMark: 1024 })
  // 异步迭代读取流的每一个数据块
  for await (const chunk of readStream) {
    // 写入返回false代表缓冲区已满,等待drain事件触发后再继续写入
    if (!writeStream.write(chunk)) {
      await new Promise(resolve => writeStream.once('drain', resolve))
    }
  }
}

调用代码修改:需要在所有文件合并完成后手动关闭写入流:

async function runConcat() {
  const fileWriter = createWriteStream('concatBins.bin', { flags: 'w' })
  let writtenLen = 0
  const fileList = {
    0: "foo.bin",
    1: "bar.bin"
  }
  for (const [key, value] of Object.entries(fileList)) {
    await concatFile(value, fileWriter)
    writtenLen = fileWriter.bytesWritten
    console.log('bytes written ' + writtenLen)
  }
  // 全部写入完成后关闭写入流,确保缓冲区数据全部刷入磁盘
  await new Promise(resolve => fileWriter.end(resolve))
}

runConcat().catch(err => console.error('文件合并失败:', err))

方案2:使用内置stream.pipeline(更省心,自动处理所有流逻辑)

Node.js内置的pipeline方法会自动处理背压、错误捕获、流销毁,无需手动监听事件:

const { pipeline } = require('stream/promises')
const { createReadStream, createWriteStream } = require('fs')

async function concatFile (filename, writeStream) {
  await pipeline(
    createReadStream(filename, { highWaterMark: 1024 }),
    writeStream,
    // 配置end为false,避免写完单个文件就自动关闭目标写入流
    { end: false }
  )
}

调用逻辑和方案1完全一致。

内容的提问来源于stack exchange,提问作者rookie

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 07:48:04