如何用Promises或Rx实现Bookshelf/Knex批量SQL插入的顺序执行?
如何用Promise(或RxJS)实现顺序执行的批量数据插入
我来帮你解决这个顺序执行Promise的问题——其实核心就是要让每一次数据库操作都等上一次完成后再启动,避免一下子发起大量并发请求导致出错。下面给你几种切实可行的方案,从最简单的async/await到RxJS的实现都有:
方法一:用async/await + 普通for循环(最直观易懂)
这是目前最推荐的方式,async/await能让异步代码看起来像同步逻辑,完美适配顺序执行的需求。如果是按1000行批次插入,我们可以先把数据分组,再逐个批次等待插入完成:
async function insertDataInBatches() { const totalRows = parser.getCount(); const batchSize = 1000; // 每批次1000行 for (let batchStart = 0; batchStart < totalRows; batchStart += batchSize) { const batchEnd = Math.min(batchStart + batchSize, totalRows); const batchRows = []; // 先收集当前批次的有效数据 for (let x = batchStart; x < batchEnd; x++) { const row = parser.getRow(x); // 遇到空行或关键字段为空就终止所有插入 if (_.isEmpty(row) || _.isEmpty(row.ch) || _.isEmpty(row.location)) { batchStart = totalRows; // 跳出外层循环 break; } const locationData = breakDownLocation(row); // 假设这是同步处理逻辑 batchRows.push(locationData); } if (batchRows.length === 0) break; // 等待当前批次插入完成,再进入下一批次 await Bookshelf.Model.forge().query().insert(batchRows); // 也可以直接用Knex:await knex('your_table_name').insert(batchRows); console.log(`完成第${Math.floor(batchStart/batchSize)+1}批次,共${batchRows.length}条数据`); } console.log('所有数据插入任务完成!'); } // 调用函数并处理可能的错误 insertDataInBatches().catch(err => { console.error('插入过程出错:', err); });
如果你的需求是逐行顺序插入(而非批次),也可以简化代码,直接在循环里await单条数据的插入:
async function insertRowsOneByOne() { const totalRows = parser.getCount(); for (let x = 0; x < totalRows; x++) { const row = parser.getRow(x); if (_.isEmpty(row) || _.isEmpty(row.ch) || _.isEmpty(row.location)) break; const locationData = breakDownLocation(row); // 等待当前行插入完成再处理下一行 await Bookshelf.Model.forge(locationData).save(); console.log(`完成第${x+1}条数据插入`); } console.log('所有数据插入完成'); } insertRowsOneByOne().catch(err => console.error('插入出错:', err));
方法二:用Promise链式调用(兼容旧环境)
如果你的运行环境不支持async/await(比如老版本Node.js),可以用Promise的链式调用模拟顺序执行:
function insertDataInOrder() { const totalRows = parser.getCount(); const batchSize = 1000; let currentPromise = Promise.resolve(); // 初始Promise作为链条起点 let batchStart = 0; while (batchStart < totalRows) { const batchEnd = Math.min(batchStart + batchSize, totalRows); // 把当前批次的插入逻辑加入Promise链 currentPromise = currentPromise.then(() => { const batchRows = []; let isValidBatch = true; for (let x = batchStart; x < batchEnd; x++) { const row = parser.getRow(x); if (_.isEmpty(row) || _.isEmpty(row.ch) || _.isEmpty(row.location)) { isValidBatch = false; break; } const locationData = breakDownLocation(row); batchRows.push(locationData); } if (!isValidBatch || batchRows.length === 0) return Promise.resolve(); // 执行当前批次插入 return Bookshelf.Model.forge().query().insert(batchRows); }).then(() => { console.log(`完成第${Math.floor(batchStart/batchSize)+1}批次插入`); batchStart += batchSize; }); } // 处理最终完成和错误 return currentPromise .then(() => console.log('所有数据插入完成')) .catch(err => console.error('插入过程出错:', err)); } insertDataInOrder();
方法三:用RxJS实现顺序执行
如果你习惯用RxJS处理异步流,concatMap操作符是关键——它会等待前一个Observable完成后,再处理下一个,完美适配顺序执行的需求:
import { from, range } from 'rxjs'; import { concatMap, takeWhile } from 'rxjs/operators'; function insertDataWithRxJS() { const totalRows = parser.getCount(); const batchSize = 1000; // 创建批次索引的Observable range(0, Math.ceil(totalRows / batchSize)).pipe( concatMap(batchIndex => { const batchStart = batchIndex * batchSize; const batchEnd = Math.min(batchStart + batchSize, totalRows); const batchRows = []; let isValid = true; for (let x = batchStart; x < batchEnd; x++) { const row = parser.getRow(x); if (_.isEmpty(row) || _.isEmpty(row.ch) || _.isEmpty(row.location)) { isValid = false; break; } const locationData = breakDownLocation(row); batchRows.push(locationData); } if (!isValid || batchRows.length === 0) { return from(Promise.resolve()); // 终止后续处理 } // 将插入Promise转为Observable return from(Bookshelf.Model.forge().query().insert(batchRows)); }), // 可选:遇到无效行就停止流 takeWhile(() => { const nextRowIndex = Math.min(batchStart + batchSize, totalRows); const nextRow = parser.getRow(nextRowIndex); return !_.isEmpty(nextRow) && !_.isEmpty(nextRow.ch) && !_.isEmpty(nextRow.location); }) ).subscribe({ next: () => console.log('一个批次插入完成'), complete: () => console.log('所有数据插入任务完成'), error: err => console.error('插入出错:', err) }); } insertDataWithRxJS();
为什么原来的代码会出错?
你原来的for循环是同步执行的,会一次性发起所有的数据库插入Promise,瞬间产生大量并发请求,导致数据库连接池耗尽、超时或者抛出错误。而上面的方案都是通过强制顺序执行,控制每次只处理一个批次(或一行),从根源上避免了并发过高的问题。
内容的提问来源于stack exchange,提问作者juliano.net
相关产品推荐
相关产品推荐

