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

Node.js中逐行读文件异步推SQS的并行执行优化问询

问题

刚接触异步/Promise概念,正在优化Node.js编写的Lambda函数以实现并行执行。需求如下:

  • 从约10万行的大文件中顺序逐行读取
  • 每行生成一条消息推送到SQS
  • 必须等所有入队操作完成后函数才结束

现有代码能运行,但有个问题:会先把所有行读完,再批量执行入队操作——因为是把所有Promise收集完才调用Promise.all。想要实现的效果是:文件读取保持顺序,但入队操作能和读取交替进行,比如读几行后,之前的入队操作完成,同时继续读后续行。预期日志输出类似:

read line, line 1
read line, line 2
Enqueued message, message 1
read line, line 3
Enqueued message, message 2
Enqueued message, message 3
...

现有代码:

exports.queueUpdates = async (filePath) => {
  return new Promise((resolve, reject) => {
    const rl = readline.createInterface({
      input: fs.createReadStream(filePath),
      crlfDelay: Infinity
    });

    var queuePromises = [];

    rl.on('line', (line) => {
      console.log("read line", line);
      var message = exports.queueMessageForLine(line); // 该函数返回要发送到SQS的JSON格式消息
      if (message !== null) {
        console.log("Pushing message", message);
        queuePromises.push(
          sqs.sendMessage(message).promise()
          .then((result) => {
            console.log("Enqueued message", message, result);
            return message;
          })
          .catch((err) => {
            console.error("Failed adding message to queue", err);
            return message;
          })
        );
      }
    }).on('close', () => {
      console.log("file read");
      Promise.all(queuePromises).then((results) => {
        resolve(results)
      })
    }).on('error', (err) => {
      console.error(err, "Error in reading the file contents");
      reject();
    });
  });
};
解决方案

你要的效果完全能实现。核心思路是:别等所有行读完才处理异步请求,而是读一行就发起对应的SQS入队请求,但得控制并发数量——毕竟10万行直接全发起请求容易触发SQS限流,也会拖垮Lambda性能。同时还能保持文件读取的顺序性。

关键修改点

  1. 不再一次性收集所有Promise,而是维护一个「正在执行的Promise池」,限制并发数量(比如同时最多10个入队请求)
  2. 读取每一行时,先检查并发池是否已满:未满就立即发起入队请求;已满则等池中有一个Promise完成后再继续
  3. 最后等待所有未完成的Promise执行完毕再结束函数

修改后的代码

const fs = require('fs');
const readline = require('readline');

exports.queueUpdates = async (filePath) => {
  return new Promise((resolve, reject) => {
    const rl = readline.createInterface({
      input: fs.createReadStream(filePath),
      crlfDelay: Infinity
    });

    const concurrencyLimit = 10; // 并发数可根据SQS限流规则调整,默认建议10-50
    let activePromises = [];
    let allPromises = [];

    rl.on('line', async (line) => {
      console.log("read line", line);
      const message = exports.queueMessageForLine(line);
      if (message === null) return;

      console.log("Pushing message", message);
      // 创建入队Promise
      const enqueuePromise = sqs.sendMessage(message).promise()
        .then((result) => {
          console.log("Enqueued message", message, result);
          return message;
        })
        .catch((err) => {
          console.error("Failed adding message to queue", err);
          return message;
        })
        .finally(() => {
          // 执行完成后从活跃池中移除
          activePromises = activePromises.filter(p => p !== enqueuePromise);
        });

      activePromises.push(enqueuePromise);
      allPromises.push(enqueuePromise);

      // 达到并发限制时,等待任意一个活跃Promise完成
      if (activePromises.length >= concurrencyLimit) {
        await Promise.race(activePromises);
      }
    })
    .on('close', async () => {
      console.log("file read");
      // 等待所有入队操作完成
      const results = await Promise.all(allPromises);
      resolve(results);
    })
    .on('error', (err) => {
      console.error(err, "Error in reading the file contents");
      reject(err);
    });
  });
};

代码说明

  • concurrencyLimit:控制同时发起的SQS请求数,避免触发SQS的请求限流(SQS默认每秒最多1000个请求,可根据实际情况调整)
  • activePromises:跟踪当前正在执行的入队请求,数量达标时用Promise.race等一个请求完成,再继续处理下一行
  • allPromises:收集所有入队Promise,最后用Promise.all确保所有消息都处理完毕才结束函数
  • 这样既保证了文件按顺序读取,又能让入队操作和读取交替进行,完全符合你预期的日志输出效果

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 04:10:14