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

如何限制fs.createReadStream同时打开的文件数量不超过指定阈值

问题根因

你的限流逻辑失效和两个核心问题有关:

  1. readline的line事件会连续同步触发,不会等待你异步回调里的await逻辑执行完成,瞬间就会有数百个回调同时执行,全部通过stallIfNeeded的检查(此时files.open还没来得及累加),之后批量打开读取流超出阈值。
  2. 逐行遍历源文件时,每行都重新打开一次待比对的大文件、全量扫描,时间复杂度是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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 13:57:03