Mongoose与Node.js事务中集成sendCashback函数的实现问题
如何在Mongoose事务中集成外部函数保证一致性?
我已经基于Mongoose和Node.js实现了事务功能且运行正常,现在需要在事务中调用外部的sendCashback函数,要求事务和该函数操作要么全部成功,要么全部回滚,但当前集成后无法正常工作,请问如何正确集成sendCashback函数以实现事务一致性?
以下是我的代码:
var transaction: any let cashback_sent: any try{ const transactiontimectrl = Number(process.env.TRANSACTION_TIME_STOP) let dt = new Date(); dt.setSeconds(dt.getSeconds() - transactiontimectrl); const timectrl = dt.toISOString(); console.log('teste transaction1') console.log(timectrl) const session = await mongoose.startSession(); await session.withTransaction(async () => { const transactionRes = await Transaction.findOne( { from: from, //to: to, //value: valuetodebit, createdAt: {$gt: timectrl} }).session(session) if(!transactionRes){ transaction = await Transaction.create([ { from, to, value: value, division_factor: destinationUserGroup.division_factor, title: title || null, description: description || null, type: type || null, hasCashback: hasCashback, realmoney: realmoney, valuetodebit: valuetodebit }],{ session }) }else{ throw new Error('Erro, uma transação semelhante foi realizada recentemente') } if (!transaction) { throw new Error('Erro, tente novamente mais tarde') } console.log('@sadihjaisvq3') let fromBalance //Saldo real if (realmoney) { fromBalance = await Balance.findOne({ type: IBalanceType.GENERAL, user: from }).session(session) } //Saldo cashback else { fromBalance = await Balance.findOne({ type: IBalanceType.CASHBACK, user: from }).session(session) } console.log('****criar balance2') let toBalance = await Balance.findOne({ type: IBalanceType.GENERAL, user: to }).session(session) console.log('****criar balance1', toBalance) if (!toBalance) { console.log('****criar balance') toBalance = await Balance.create( { user: to }) } let toLivetPay = await Balance.findOne({ type: IBalanceType.GENERAL, user: toLivet }).session(session) if (!fromBalance || !toBalance || !toLivetPay) throw new Error() if (!credit) { fromBalance.value = fromBalance.value + parseFloat(valuetodebit) * -1 } // fromBalance.value = fromBalance.value + parseFloat(valuetodebit) * -1 toBalance.value = toBalance.value + parseFloat(value) let valorCash = Number(value) / 0.75 let valorLivet = Number(valorCash) * 0.25 toLivetPay.value = toLivetPay.value + parseFloat(String(valorLivet)) await fromBalance.save() await toBalance.save() await toLivetPay.save() if (hasCashback) { if (!saque) { //EXTERNAL FUNCTION--------------------------------------------------------- cashback_sent = await sendCashback( transaction[0], body.originUser!, body.destinationUser! ) //------------------------------------------------------------------------ } if (!cashback_sent) { throw new Error('Erro ao distribuir cashback') } console.log('depois do cashback') } }) } catch (err) { console.error(err); throw err; }
解决方案
1. 核心问题说明
MongoDB事务仅能管控自身的数据库操作,外部函数(比如sendCashback)不在事务覆盖范围内。如果在事务内部调用外部函数,一旦后续操作失败导致事务回滚,已执行的外部操作无法撤销,必然出现数据不一致。
2. 调整执行逻辑
正确做法是先确保Mongo事务全部提交成功,再执行外部函数,同时通过状态标记和补偿机制处理外部函数执行失败的场景:
步骤1:给Transaction文档新增状态字段
添加cashbackStatus字段标记cashback执行状态,可选值:
pending:事务已提交,待执行cashbackcompleted:cashback执行成功failed:cashback执行失败
修改事务内创建Transaction的代码:
transaction = await Transaction.create([ { from, to, value: value, division_factor: destinationUserGroup.division_factor, title: title || null, description: description || null, type: type || null, hasCashback: hasCashback, realmoney: realmoney, valuetodebit: valuetodebit, cashbackStatus: hasCashback && !saque ? 'pending' : 'completed' // 新增状态字段 }],{ session })
步骤2:将外部函数调用移到事务外部
事务提交成功后再执行sendCashback,并更新状态:
const session = await mongoose.startSession(); try { await session.withTransaction(async () => { // ... 保留原事务内的所有数据库操作,移除sendCashback调用 }); // 事务提交成功后处理cashback if (hasCashback && !saque) { try { cashback_sent = await sendCashback( transaction[0], body.originUser!, body.destinationUser! ); if (cashback_sent) { await Transaction.findByIdAndUpdate(transaction[0]._id, { cashbackStatus: 'completed' }); } else { await Transaction.findByIdAndUpdate(transaction[0]._id, { cashbackStatus: 'failed' }); throw new Error('Erro ao distribuir cashback'); } } catch (err) { await Transaction.findByIdAndUpdate(transaction[0]._id, { cashbackStatus: 'failed' }); throw err; } } } catch (err) { console.error(err); // 事务失败时无需处理cashback throw err; } finally { session.endSession(); // 确保session被正确关闭 }
3. 额外保障措施
- 幂等性处理:给
sendCashback添加基于Transaction ID的唯一标识校验,避免同一事务重复执行cashback。 - 定时重试机制:实现定时任务,扫描
cashbackStatus为pending或failed的记录,自动重试sendCashback,解决临时网络/服务不可用问题。 - 错误日志与告警:记录cashback失败的详细日志并设置告警,便于人工介入处理极端情况。
内容的提问来源于stack exchange,提问作者Maiko Souza
相关产品推荐
相关产品推荐

