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性能。同时还能保持文件读取的顺序性。
关键修改点
- 不再一次性收集所有Promise,而是维护一个「正在执行的Promise池」,限制并发数量(比如同时最多10个入队请求)
- 读取每一行时,先检查并发池是否已满:未满就立即发起入队请求;已满则等池中有一个Promise完成后再继续
- 最后等待所有未完成的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
相关产品推荐
相关产品推荐

