Connector/Node.js中collection.add().execute()仅首次执行生效问题求助
MySQL Document Store批量插入仅首次成功问题排查与修复
问题描述
我正在测试MySQL Document Store,为了和现有关系表做合理对比,尝试把一张约32万条记录的表转换成集合。我用Connector/Node.js实现这个操作,为了避免性能问题,每次尝试插入1万条记录,但只有首次插入的1万条成功入库,后续批次都没有数据插入,而且已经确认首次数据结构正确,不存在覆盖情况。
原代码
const mysqlx = require('@mysql/xdevapi'); const config = { password: 'notMyPassword', user: 'notMyUser', host: 'notMyHost', port: 33060, schema: 'sample' }; var mySchema; var myCollection; var recCollection = []; mysqlx.getSession(config).then(session => { mySchema = session.getSchema('sample'); mySchema.dropCollection('sample_test'); mySchema.createCollection('sample_test'); myCollection = mySchema.getCollection('sample_test'); var myTable = mySchema.getTable('sampledata'); return myTable.select('FormDataId','FormId','DateAdded','DateUpdated','Version','JSON').orderBy('FormDataId').execute(); }).then(result => { console.log('we have a result to analyze...'); var tmp = result.fetchOne(); while(tmp !== null && tmp !== '' && tmp !== undefined){ var r = tmp; var myRecord = { 'dateAdded': r[2], 'dateUpdated': r[3], 'version': r[4], 'formId': r[1], 'dataId': r[0], 'data': r[5] }; recCollection.push(myRecord); if (recCollection.length >= 10000){ console.log('inserting 10000'); try { myCollection.add(recCollection).execute(); } catch(ex){ console.log('error: ' + ex); } recCollection.length = 0; } tmp = result.fetchOne(); } });
问题根源
- 异步操作未等待:
myCollection.add().execute()是异步Promise操作,但原代码没有等待它完成就立即清空了recCollection数组,导致后续插入任务拿到的是空数组,无法写入数据。 - 集合操作未同步:
dropCollection和createCollection也是异步操作,未等待完成就直接进行后续操作,可能导致集合状态异常。 - 剩余记录未处理:循环结束后,数组中剩余的不足1万条记录没有被插入。
修复后的代码
const mysqlx = require('@mysql/xdevapi'); const config = { password: 'notMyPassword', user: 'notMyUser', host: 'notMyHost', port: 33060, schema: 'sample' }; mysqlx.getSession(config).then(async session => { const mySchema = session.getSchema('sample'); // 忽略集合不存在的错误,等待删除完成 await mySchema.dropCollection('sample_test').catch(() => {}); // 等待集合创建完成 const myCollection = await mySchema.createCollection('sample_test'); const myTable = mySchema.getTable('sampledata'); const result = await myTable.select('FormDataId','FormId','DateAdded','DateUpdated','Version','JSON') .orderBy('FormDataId').execute(); console.log('开始处理数据...'); let recCollection = []; let tmp; while((tmp = result.fetchOne()) !== null){ const r = tmp; const myRecord = { dateAdded: r[2], dateUpdated: r[3], version: r[4], formId: r[1], dataId: r[0], data: r[5] }; recCollection.push(myRecord); if (recCollection.length >= 10000){ console.log('插入10000条记录'); try { // 等待插入操作完成后再继续 await myCollection.add(recCollection).execute(); } catch(ex){ console.error('插入失败:', ex); } recCollection = []; } } // 处理循环结束后剩余的记录 if(recCollection.length > 0){ console.log(`插入剩余${recCollection.length}条记录`); try { await myCollection.add(recCollection).execute(); } catch(ex){ console.error('剩余记录插入失败:', ex); } } console.log('数据迁移完成'); await session.close(); }).catch(err => { console.error('会话连接或操作失败:', err); });
关键修复说明
- 改用
async/await处理所有异步操作,确保每一步操作完成后再执行后续逻辑,避免异步竞态问题。 - 等待集合删除、创建操作完成,保证集合状态正常。
- 插入操作必须等待完成后再清空数组,确保数据被正确写入。
- 增加剩余记录的处理逻辑,避免数据遗漏。
- 优化错误输出,方便排查问题。
- 最后关闭数据库会话,释放资源。
内容的提问来源于stack exchange,提问作者Loco
相关产品推荐
相关产品推荐

