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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 23:45:41