如何在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
相关产品推荐
相关产品推荐

