You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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)
        }
      })
  }
}

关键修改点说明:

  1. 连贯的Observable流:用flatMap把两个插入操作和事务提交串起来,确保前一个操作完成后才执行下一个,错误会自动传播到后续的recoverWith。
  2. 统一错误处理:在事务操作链的末尾添加recoverWith,捕获任何异常后立即调用abortTransaction回滚事务,然后把错误继续抛出。
  3. 重构事务入口:updateEmployeeInfo现在接收ClientSession并返回SingleObservable[Completed],把事务的整个生命周期(开启、执行操作、提交/回滚)封装在一个流里,避免异步操作的时序问题。
  4. 避免直接使用subscribe:subscribe会脱离当前的Observable流,导致错误无法被事务逻辑捕获,改用flatMap来组合操作,让错误能被正确处理。

这样修改后,只要事务内的任何一步操作失败,整个事务都会被回滚,第一个插入的文档不会被持久化。

内容的提问来源于stack exchange,提问作者Teimuraz

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.12 05:37:36