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
相关产品推荐
相关产品推荐

