Node.js pipeline()循环合并文件块仅执行首次迭代问题求助
问题根源
stream/promises的pipeline函数默认会在执行完成后自动销毁所有传入的流,包括你创建的目标写入流writeStream。第一次循环执行await pipeline(...)后,写流已经被关闭并销毁,后续循环再使用这个流时,它处于不可写入状态,因此没有任何输出;同时流关闭后的写入操作通常不会抛出明显错误,导致控制台无报错信息。
修复方案
推荐两种可靠的修复方式,可根据场景选择:
方案1:保持单写流,阻止pipeline自动关闭它
修改pipeline调用,添加end: false选项,让pipeline完成后不关闭写流,直到所有块处理完毕再手动关闭。同时在错误处理中确保销毁写流,避免资源泄漏:
import { pipeline } from "stream/promises"; import fs from "fs"; import path from "path"; async function combineChunks( tempFolder, numberedFiles, outputPath, combinedFileSize ) { const writeStream = fs.createWriteStream(outputPath, { flags: "a" }); try { for (const file of numberedFiles) { const filePath = path.join(tempFolder, file); if (!filePath.startsWith(tempFolder)) { throw new Error("Invalid file path detected"); } // 关键修改:添加end: false,阻止pipeline关闭写流 await pipeline( fs.createReadStream(filePath), writeStream, { end: false } ); await fs.promises.unlink(filePath); } // 所有块写入完成后,手动关闭写流并等待写入完成 writeStream.end(); await new Promise((resolve, reject) => { writeStream.on("finish", resolve); writeStream.on("error", reject); }); // 异步获取文件信息,避免阻塞事件循环 const stats = await fs.promises.stat(outputPath); const fileSize = stats.size; const allowedDocumentExtensions = /\.(jpg|jpeg|png|gif|bmp|svg|webp|txt|doc|docx|odt|xls|xlsx|ods|ppt|pptx|odp|pdf)$/i; if (fileSize !== combinedFileSize) { await fs.promises.unlink(outputPath); throw new Error("File size mismatch"); } if (!allowedDocumentExtensions.test(outputPath)) { await fs.promises.unlink(outputPath); throw new Error("Inappropriate file type detected"); } return true; } catch (error) { // 出错时销毁写流,释放资源 writeStream.destroy(error); console.error("Error combining files:", error); return false; } }
方案2:每次循环创建新的追加写流
如果担心单流管理的复杂度,可以把写流的创建移到循环内部,每次处理一个块时打开追加模式的写流,pipeline完成后自动关闭它:
import { pipeline } from "stream/promises"; import fs from "fs"; import path from "path"; async function combineChunks( tempFolder, numberedFiles, outputPath, combinedFileSize ) { try { for (const file of numberedFiles) { const filePath = path.join(tempFolder, file); if (!filePath.startsWith(tempFolder)) { throw new Error("Invalid file path detected"); } // 每次循环创建新的追加写流 const writeStream = fs.createWriteStream(outputPath, { flags: "a" }); await pipeline( fs.createReadStream(filePath), writeStream ); await fs.promises.unlink(filePath); } // 异步获取文件信息 const stats = await fs.promises.stat(outputPath); const fileSize = stats.size; const allowedDocumentExtensions = /\.(jpg|jpeg|png|gif|bmp|svg|webp|txt|doc|docx|odt|xls|xlsx|ods|ppt|pptx|odp|pdf)$/i; if (fileSize !== combinedFileSize) { await fs.promises.unlink(outputPath); throw new Error("File size mismatch"); } if (!allowedDocumentExtensions.test(outputPath)) { await fs.promises.unlink(outputPath); throw new Error("Inappropriate file type detected"); } return true; } catch (error) { console.error("Error combining files:", error); return false; } }
额外优化建议
- 替换同步的
fs.statSync为异步的await fs.promises.stat,避免阻塞事件循环。 - 手动关闭写流后监听
finish事件,确保所有数据写入磁盘后再进行文件验证,避免读取到不完整的文件。
内容的提问来源于stack exchange,提问作者Sam Caut
相关产品推荐
相关产品推荐

