Node.js循环并行请求过慢,寻求高效优化方案
问题描述
我正在开发一项功能,需要处理超过10k条数据的对象,逐条取出数据发送至API处理,获取响应后再处理下一条。目前用async库的mapLimit方法实现,功能正常但数据量太大导致遍历耗时过长。这个功能需要持续运行,遍历完所有数据后等待几秒再重复执行。我试过用child_process创建子进程,把数据分片分配给各进程处理,功能正常但内存消耗极高,导致服务器上其他进程因内存不足故障,就算进程退出后问题还存在。请问怎么实现更快的处理速度?
现有代码
获取钱包列表的代码
getListofWallet = async () => { try { const USDT = await usdt.query(sql` SELECT * FROM USDT ORDER BY id DESC; `); let counter = 0; let completion = 0; async.mapLimit(USDT, 6, async (user) => { let userDetail = { email: user.email, id: user.user_id, address: user.address } try { await this.getWalletTransaction(userDetail); completion++; } catch (TronGridException) { completion++; console.log(":: A TronGrid Exception Occured"); console.log(TronGridException); } if (USDT.length == completion || USDT.length == (completion-5)) { setTimeout(() => { this.getListofWallet(); }, 60000); console.log('~~~~~~~ Finished Wallet Batch ~~~~~~~~~'); } }); } catch (error) { console.log(error); console.log('~~~~~~~Restarting TronWallet File after Crash ~~~~~~~~~'); this.getListofWallet(); } }
处理数据并执行必要操作的代码
getWalletTransaction = async (walletDetail) => { const config = { headers: { 'TRON-PRO-API-KEY': process.env.api_key, 'Content-Type': 'application/json' } }; const getTransactionFromAddress = await axios.get(`https://api.trongrid.io/v1/accounts/${walletDetail.address}/transactions/trc20`, config); const response = getTransactionFromAddress.data; const currentTimeInMillisecond = 1642668127000; //1632409548000 response.data.forEach(async (value) => { if (value.block_timestamp >= currentTimeInMillisecond && value.token_info.address == "TR7NHqjeKQxGTCi8q8ZY4pL8otSzgjLj6t") { let doesHashExist = await transactionCollection.query(sql`SELECT * FROM transaction_collection WHERE txhash=${value.transaction_id};`); if (doesHashExist.length == 0) { if (walletDetail.address == value.to) { const checkExistence = await CryptoTransactions2.query(sql` SELECT * FROM CryptoTransaction2 WHERE txHash=${value.transaction_id}; `); if (checkExistence.length == 0) { const xCollection = { collection: "CryptoTransaction2", queryObject: { currency: "USDT", sender: ObjectID("60358d21ec2b4b33e2fcd62e"), receiver: ObjectID(walletDetail.id), amount: parseFloat(tronWeb.fromSun(value.value)), txHash: value.transaction_id, description: "New wallet Deposit " + "60358d21ec2b4b33e2fcd62e" + " into " + value.to, category: "deposit", createdAt: new Date(), updatedAt: new Date(), }, }; await new MongoDbService().createTransaction(xCollection); //create record inside cryptotransactions db. await CryptoTransactions2.query(sql`INSERT INTO CryptoTransaction2 (txHash) VALUES (${value.transaction_id})`) } } } } }); }
优化方案
1. 合理调整并发数,适配API限流
- 先确认TronGrid API的官方限流规则(比如允许的并发请求数、请求频率),在不触发429错误的前提下,适当提高
async.mapLimit的并发数(比如从6调整到10-20),提升处理速度。 - 给axios添加重试机制,比如使用
axios-retry库,遇到网络波动或429限流错误时自动重试,减少因请求失败导致的等待时间。
2. 优化数据库操作,减少IO开销
- 批量查询替代单条查询:在
getWalletTransaction中,不要每条交易都单独查询数据库,先收集当前钱包下所有符合条件的transaction_id,再一次性批量查询是否存在,减少数据库连接次数。 - 添加索引:给
transaction_collection.txhash和CryptoTransactions2.txHash字段添加唯一索引,大幅提升查询速度,同时避免重复插入。 - **避免SELECT ***:只查询需要的字段(比如只查
txhash),减少数据传输量和内存占用。
3. 修复内存泄漏与任务调度逻辑
- 修正任务结束判断:当前代码通过
completion变量判断批次结束的逻辑存在误差,且可能导致重复启动新批次。改为等待async.mapLimit全部完成后再启动下一批:// 替换原有mapLimit调用 await async.mapLimit(USDT, 10, async (user) => { // ... 原有处理逻辑 }); console.log('~~~~~~~ Finished Wallet Batch ~~~~~~~~~'); setTimeout(() => { this.getListofWallet(); }, 60000); - 排查子进程内存问题:如果继续用子进程,确保子进程退出时彻底关闭所有数据库连接、清除事件监听;用
node --inspect工具分析内存快照,定位泄漏点。
4. 增量处理数据,减少重复计算
- 不要每次全量查询
USDT表,记录上次处理的最大id或最后处理时间,下次只查询新增的数据:
这样每次只处理新增钱包,避免重复处理已处理过的数据。SELECT * FROM USDT WHERE id > ${lastProcessedId} ORDER BY id DESC;
5. 批量执行数据库插入
- 把需要插入
CryptoTransactions2的记录收集成数组,最后一次性批量插入,减少数据库IO次数:let insertBatch = []; // 遍历交易时收集需要插入的txHash if (checkExistence.length == 0) { insertBatch.push(value.transaction_id); // ... 其他MongoDB插入逻辑 } // 批量插入 if (insertBatch.length > 0) { await CryptoTransactions2.query(sql`INSERT INTO CryptoTransaction2 (txHash) VALUES ${insertBatch.map(h => `(${h})`).join(',')}`); }
6. 替换async库为原生实现(可选)
- 用原生
Promise.allSettled结合手动分批,替代async.mapLimit,更灵活控制并发:const batchSize = 10; for (let i = 0; i < USDT.length; i += batchSize) { const batch = USDT.slice(i, i + batchSize); await Promise.allSettled(batch.map(user => this.getWalletTransaction({ email: user.email, id: user.user_id, address: user.address }))); }
内容的提问来源于stack exchange,提问作者vhiktordom
相关产品推荐
相关产品推荐

