如何限制fs.createReadStream同时打开的文件数量不超过指定阈值
问题根因
你的限流逻辑失效和两个核心问题有关:
readline的line事件会连续同步触发,不会等待你异步回调里的await逻辑执行完成,瞬间就会有数百个回调同时执行,全部通过stallIfNeeded的检查(此时files.open还没来得及累加),之后批量打开读取流超出阈值。- 逐行遍历源文件时,每行都重新打开一次待比对的大文件、全量扫描,时间复杂度是O(n*m),n和m都是20万行的量级,执行效率极低,完全没必要。
最优解决方案(推荐)
20万行文本的内存占用极低,完全可以先把待比对文件的所有行预读到内存的Set结构中,再逐行遍历源文件做判断,全程只需要打开3个文件流(两个读、一个写),彻底解决流数量超限的问题,同时执行效率提升上百倍。
示例代码如下:
var fs = require('fs'); const readline = require('readline'); // 预加载待比对文件的所有行到Set async function loadTargetLines() { const targetSet = new Set(); const readStream = fs.createReadStream("./largeFileOfUncertianValues.csv"); const rl = readline.createInterface({ input: readStream, crlfDelay: Infinity }); return new Promise(res => { rl.on("line", line => { // 这里可以按需做format处理,和后面的源文件行格式对齐 targetSet.add(line.trim()); }); rl.on("close", () => res(targetSet)); }); } (async () => { // 先加载待比对的行 const targetLines = await loadTargetLines(); console.log("待比对文件加载完成,共" + targetLines.size + "行"); const sourceStream = fs.createReadStream('largeFileOfKnownValues.csv'); const sourceRl = readline.createInterface({ input: sourceStream, crlfDelay: Infinity }); const writeStream = fs.createWriteStream("missing.csv"); let cntr = 0; await new Promise(sRes => { sourceRl.on("line", line => { const lineFormatted = formatLine(line); const skipLine = checkLine(lineFormatted); if (skipLine) return; const exists = targetLines.has(lineFormatted.trim()); if (!exists) { writeStream.write(line + "\n"); } if (++cntr % 50 == 0) { console.log("another 50 done"); console.log(cntr + " so far"); } }); sourceRl.on("close", sRes); }); writeStream.close(); console.log("处理完成,共" + cntr + "行"); })();
若必须保留每行开流的逻辑的修复方案
如果有特殊场景不能预加载全量数据,需要保留原有的逐行开流扫描的逻辑,只需要对源readline做并发控制即可,避免一次性触发太多line回调:
// 其他原有逻辑保持不变,修改主流程的line事件处理 (async () => { const sourceStream = fs.createReadStream('largeFileOfKnownValues.csv'); const sourceRl = readline.createInterface({ input: sourceStream, crlfDelay: Infinity, highWaterMark: 1024 * 1024 // 调整缓冲区大小,避免一次性加载过多行 }); let writeStream = fs.createWriteStream("missing.csv") let cntr = 0; // 并发数控制,设为你需要的阈值,比如50 const CONCURRENT_LIMIT = 50; let activeCount = 0; // 待处理的行队列 const queue = []; // 处理队列中的行 async function processQueue() { if (activeCount >= CONCURRENT_LIMIT || queue.length === 0) return; activeCount++; const line = queue.shift(); const lineFormatted = formatLine(line); const skipLine = checkLine(lineFormatted); if (!skipLine) { const exists = await checkLineExists(lineFormatted); if (!exists) writeStream.write(line + "\n"); } if (++cntr % 50 == 0) { console.log("another 50 done"); console.log(cntr + " so far"); } activeCount--; // 继续处理下一个 processQueue(); } await new Promise(sRes => { sourceRl.on("line", line => { queue.push(line); processQueue(); // 队列积压过多时暂停读取 if (queue.length > CONCURRENT_LIMIT * 2) { sourceRl.pause(); } }); sourceRl.on("close", async () => { // 等待队列全部处理完成 while (queue.length > 0 || activeCount > 0) { await new Promise(r => setTimeout(r, 100)); } sRes(); }); }); writeStream.close(); })();
内容的提问来源于stack exchange,提问作者async await
相关产品推荐
相关产品推荐

