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

Mongo与Oracle数据桥接脚本:并行查询Mongo遇异步执行问题求助

解决Mongo查询并行执行后等待所有结果再写入Oracle的问题

你遇到的核心问题是:原来的Mongo查询用的是回调式异步API,没有把这些操作包装成可被async/await等待的Promise,导致脚本还没等三个集合的查询和数据处理完成,就直接跳到Oracle写入步骤了,自然拿不到正确的处理后数据。

下面是具体的修复方案,一步步帮你理顺异步流程:

1. 把回调式Mongo查询重构为返回Promise的函数

MongoDB的find().exec()可以直接返回Promise(不传回调参数即可),我们把RequestOne、RequestTwo改成异步函数,这样就能用await等待它们完成:

// 重构RequestOne为异步函数
async function RequestOne(dateStart, dateEnd, RowToInsert) {
  try {
    // 直接await exec()返回的Promise,拿到查询结果
    const arrAqsOneHourMoy = await CollectionOne.find({
      AQS_REF_TIME_EVENT_MSR: { $gte: dateStart, $lte: dateEnd }
    }).exec();

    if (arrAqsOneHourMoy && arrAqsOneHourMoy.length > 0) {
      let attributeOneSum = 0;
      let attributeTwoSum = 0;
      // 用let声明循环变量i,避免全局污染
      for (let i = 0; i < arrAqsOneHourMoy.length; i++) {
        attributeOneSum += arrAqsOneHourMoy[i].AttributeOne;
        attributeTwoSum += arrAqsOneHourMoy[i].AttributeTwo;
      }
      // 修正原代码的赋值错误:原来把AttributeOne赋值了两次
      RowToInsert.AttributeOne = attributeOneSum / arrAqsOneHourMoy.length;
      RowToInsert.AttributeTwo = attributeTwoSum / arrAqsOneHourMoy.length;
    }
  } catch (err) {
    console.error('RequestOne执行出错:', err);
    throw err; // 抛出错误让上层处理
  }
}

// RequestTwo用同样逻辑重构
async function RequestTwo(dateStart, dateEnd, RowToInsert) {
  try {
    const arrAqsTwoHourMoy = await CollectionTwo.find({
      AQS_REF_TIME_EVENT_MSR: { $gte: dateStart, $lte: dateEnd }
    }).exec();

    if (arrAqsTwoHourMoy && arrAqsTwoHourMoy.length > 0) {
      let attributeThreeSum = 0;
      let attributeFourSum = 0;
      for (let i = 0; i < arrAqsTwoHourMoy.length; i++) {
        attributeThreeSum += arrAqsTwoHourMoy[i].AttributeThree;
        attributeFourSum += arrAqsTwoHourMoy[i].AttributeFour;
      }
      RowToInsert.AttributeThree = attributeThreeSum / arrAqsTwoHourMoy.length;
      RowToInsert.AttributeFour = attributeFourSum / arrAqsTwoHourMoy.length;
    }
  } catch (err) {
    console.error('RequestTwo执行出错:', err);
    throw err;
  }
}

2. 修改全局函数,用async/await + Promise.all并行等待所有请求

现在RequestOne和RequestTwo都是异步函数了,我们用Promise.all并行执行所有查询,等全部完成后再执行Oracle写入:

// 给GenericFunction加上async关键字,使其成为异步函数
async function GenericFunction(time, sensorID, TEST_oracle_save) {
  const d = new Date(2018, 00, 03, 11, 00, 00, 000);
  const dateStart = d.toISOString();
  const d2 = new Date(2018, 00, 03, 11, 59, 59, 999);
  const dateEnd = d2.toISOString();

  try {
    // 用Promise.all并行执行所有Mongo查询,等待全部完成
    await Promise.all([
      RequestOne(dateStart, dateEnd, TEST_oracle_save),
      RequestTwo(dateStart, dateEnd, TEST_oracle_save)
      // 第三个请求直接加在这里:RequestThree(...)
    ]);

    // 所有查询和数据处理完成后,再连接Oracle持久化数据
    console.log('所有Mongo数据处理完成,开始写入Oracle');
    // 你的Oracle连接、写入逻辑放在这里
    // 示例:
    // const oracleConn = await oracle.connect(yourOracleConfig);
    // await oracleConn.execute('INSERT INTO your_table (...) VALUES (...)', [TEST_oracle_save.AttributeOne, ...]);
    // await oracleConn.close();

  } catch (err) {
    console.error('流程执行出错:', err);
    // 这里可以添加错误处理,比如Oracle回滚、告警等
  }
}

关键细节说明

  • Promise.all会并行执行所有传入的Promise,只有当所有Promise都成功完成后才会继续执行;如果有任何一个失败,会立即抛出错误。
  • 所有使用await的函数必须加上async关键字,否则会报错。
  • 修正了原代码中变量未声明(如i)、重复赋值的问题,避免全局变量污染和逻辑错误。
  • 用try/catch包裹异步操作,方便捕获和排查错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:27:54