Node.js中MongoDB事务报错:提交后无法调用abortTransaction
问题原因分析
报错Cannot call abortTransaction after calling commitTransaction的核心原因是异步函数未正确使用await,导致事务流程逻辑混乱:
runTransationWithRetry调用异步事务函数txnFunc时未加await,同步执行循环逻辑会导致事务操作未完成就触发下一次循环,可能重复执行事务提交/回滚操作。insertRankings和updateRankings中调用commitWithRetry时未加await,提交操作异步执行的情况下函数提前退出,后续重试逻辑可能对已提交的事务执行回滚。updateOne调用参数存在问题:直接传入整个文档会将其作为查询条件,正确写法应为指定查询条件+更新内容,否则可能出现更新不符合预期的情况。
修复方案
以下是修正后的代码,重点修复异步流程和事务操作的问题:
export async function rankingInsertOrUpdate( storedAccuracy: WithId<NFTCollectionRanking> | null, rankingDocument: NFTCollectionRanking, sortedRankingDocument: NFTSortedRanking ) { const session = client.startSession(); try { if (storedAccuracy === null) { await runTransationWithRetry(insertRankings, session); } else if (rankingDocument.accuracy > storedAccuracy.accuracy) { await runTransationWithRetry(updateRankings, session); } } catch (error: any) { throw new Error(error.message || error); } finally { await session.endSession(); } async function insertRankings(session: ClientSession) { await session.startTransaction({ readConcern: { level: 'snapshot' }, writeConcern: { w: 'majority' }, }); try { await rankings.insertOne(rankingDocument, { session }); await sortedRankings.insertOne(sortedRankingDocument, { session }); await commitWithRetry(session); } catch (error) { console.log('Caught exception during insert transaction, aborting.'); await session.abortTransaction(); throw error; } } async function updateRankings(session: ClientSession) { await session.startTransaction({ readConcern: { level: 'snapshot' }, writeConcern: { w: 'majority' }, }); try { await rankings.updateOne( { _id: rankingDocument._id }, { $set: rankingDocument }, { session } ); await sortedRankings.updateOne( { _id: sortedRankingDocument._id }, { $set: sortedRankingDocument }, { session } ); await commitWithRetry(session); } catch (error) { console.log('Caught exception during update transaction, aborting.'); await session.abortTransaction(); throw error; } } } async function runTransationWithRetry(txnFunc: (session: ClientSession) => Promise<void>, session: ClientSession) { while (true) { try { await txnFunc(session); break; } catch (error: any) { if (error.errorLabels?.includes('TransientTransactionError')) { console.log('TransientTransactionError, retrying transaction ...'); continue; } else { throw error; } } } } async function commitWithRetry(session: ClientSession) { while (true) { try { await session.commitTransaction(); console.log('Transaction committed'); break; } catch (error: any) { if (error.errorLabels?.includes('UnknownTransactionCommitResult')) { console.log('UnknownTransactionCommitResult, retrying commit operation ...'); continue; } else { console.log('Error during commit ...'); throw error; } } } }
关键修复点说明
- 所有异步函数调用(事务执行、提交、回滚、session操作)都添加
await,确保流程按顺序执行,避免事务状态混乱。 - 数据库操作(insertOne/updateOne)明确指定
session参数,确保操作在事务会话中执行。 - 修正
updateOne的参数逻辑,使用查询条件+更新内容的正确写法,避免不符合预期的更新。 - 优化错误信息输出,保留原始错误的同时提升可读性。
内容的提问来源于stack exchange,提问作者Collecto
相关产品推荐
相关产品推荐

