Scala MongoDB事务执行失败时如何实现全局回滚?
解决Scala MongoDB事务中操作失败时的回滚问题
你的问题核心在于异步操作的错误没有被正确捕获和处理——你用subscribe执行插入操作时,代码会直接继续往后走,不会等待操作完成,也不会把错误传递到事务的执行流中,导致即使第二个插入失败,事务仍然会尝试提交,最终第一个插入被持久化。
要实现任意操作失败时回滚整个事务,你需要把事务内的所有操作整合成一个连贯的Observable流,让错误能沿着流传播,从而触发事务回滚。下面是修正后的完整代码:
import org.mongodb.scala._ import org.mongodb.scala.bson.Document import scala.concurrent.Await import scala.concurrent.duration.Duration object Application extends App { val mongoClient: MongoClient = MongoClient("mongodb://localhost:27018") val database = mongoClient.getDatabase("hr") val employeesCollection = database.getCollection("employees") // Implicit functions that execute the Observable and return the results val waitDuration = Duration(5, "seconds") implicit class ObservableExecutor[T](observable: Observable[T]) { def execute(): Seq[T] = Await.result(observable.toFuture(), waitDuration) } implicit class SingleObservableExecutor[T](observable: SingleObservable[T]) { def execute(): T = Await.result(observable.toFuture(), waitDuration) } updateEmployeeInfoWithRetry(mongoClient).execute() Thread.sleep(3000) /// ------------------------- def updateEmployeeInfo(database: MongoDatabase, clientSession: ClientSession): SingleObservable[Completed] = { val eventsCollection = database.getCollection("events") val transactionOptions = TransactionOptions.builder() .readConcern(ReadConcern.SNAPSHOT) .writeConcern(WriteConcern.MAJORITY) .build() // 开启事务后,链式执行两个插入操作 clientSession.startTransaction(transactionOptions) eventsCollection.insertOne(clientSession, Document("_id" -> "123", "employee" -> 3, "status" -> Document("new" -> "Inactive", "old" -> "Active"))) .flatMap(_ => { // 第一个插入成功后执行第二个,这里会因为ID重复失败 eventsCollection.insertOne(clientSession, Document("_id" -> "123", "employee" -> 3, "status" -> Document("new" -> "Inactive", "old" -> "Active"))) }) .flatMap(_ => { // 所有操作成功后提交事务 clientSession.commitTransaction() }) .recoverWith { // 任何步骤失败时,回滚事务并抛出错误 case e: Throwable => println(s"Operation failed, aborting transaction: ${e.getMessage}") clientSession.abortTransaction().flatMap(_ => SingleObservable.failed(e)) } } def commitAndRetry(observable: SingleObservable[Completed]): SingleObservable[Completed] = { observable.recoverWith({ case e: MongoException if e.hasErrorLabel(MongoException.UNKNOWN_TRANSACTION_COMMIT_RESULT_LABEL) => { println("UnknownTransactionCommitResult, retrying commit operation ...") commitAndRetry(observable) } case e: Exception => { println(s"Exception during commit ...: $e") SingleObservable.failed(e) } }) } def runTransactionAndRetry(client: MongoClient): SingleObservable[Completed] = { client.startSession() .flatMap(clientSession => updateEmployeeInfo(client.getDatabase("hr"), clientSession)) .recoverWith({ case e: MongoException if e.hasErrorLabel(MongoException.TRANSIENT_TRANSACTION_ERROR_LABEL) => { println("TransientTransactionError, retrying transaction ...") runTransactionAndRetry(client) } }) } }
关键修改点说明:
- 连贯的Observable流:用
flatMap把两个插入操作和事务提交串起来,确保前一个操作完成后才执行下一个,错误会自动传播到后续的recoverWith。 - 统一错误处理:在事务操作链的末尾添加
recoverWith,捕获任何异常后立即调用abortTransaction回滚事务,然后把错误继续抛出。 - 重构事务入口:
updateEmployeeInfo现在接收ClientSession并返回SingleObservable[Completed],把事务的整个生命周期(开启、执行操作、提交/回滚)封装在一个流里,避免异步操作的时序问题。 - 避免直接使用
subscribe:subscribe会脱离当前的Observable流,导致错误无法被事务逻辑捕获,改用flatMap来组合操作,让错误能被正确处理。
这样修改后,只要事务内的任何一步操作失败,整个事务都会被回滚,第一个插入的文档不会被持久化。
内容的提问来源于stack exchange,提问作者Teimuraz
相关产品推荐
相关产品推荐

