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

判断MongoDB连接可用性并重建连接:寻求Scala地道实现方案

如何用地道的Scala方式处理MongoDB失效连接的重建?

我正在优化一段用于MongoDB连接与文档插入的Scala代码,原代码在连接失效(比如服务器离线)时无法自动重建连接。我做了初步改进,但希望了解更符合Scala风格的实现方式。

原代码问题

原代码中,MongoClient、Database和Collection都是初始化后固定的实例,一旦连接失效,无法重新创建:

import com.typesafe.scalalogging.LazyLogging
import org.mongodb.scala.result.InsertOneResult
import org.mongodb.scala.{Document, MongoClient, MongoCollection, MongoDatabase, Observer, SingleObservable}
import play.api.libs.json.JsResult.Exception

object MongoFactory extends LazyLogging {
  val uri: String = "mongodb+srv://*********"
  val client: MongoClient = MongoClient(uri)
  val db: MongoDatabase = client.getDatabase("db")
  val collection: MongoCollection[Document] = db.getCollection("col")

  def insertDocument(document: Document) = {
    val singleObservable: SingleObservable[InsertOneResult] = collection.insertOne(document)
    singleObservable.subscribe(new Observer[InsertOneResult] {
      override def onNext(result: InsertOneResult): Unit = println(s"onNext: $result")
      override def onError(e: Throwable): Unit = println(s"onError: $e")
      override def onComplete(): Unit = println("onComplete")
    })
  }
}

我的初步改进

我尝试将实例改为可变变量,并在插入前检查连接状态,失效则重建:

object MongoFactory extends LazyLogging {
  val uri: String = "mongodb+srv://*********"
  var client: MongoClient = MongoClient(uri)
  var db: MongoDatabase = client.getDatabase("db")
  var collection: MongoCollection[Document] = db.getCollection("col")

  def isDbDown() : Boolean = {
    try {
      client.getDatabase("db")
      false
    } catch {
      case e: Exception => true
    }
  }

  def insertDocument(document: Document) = {
    if(isDbDown()) {
      client = MongoClient(uri)
      db = client.getDatabase("db")
      collection = db.getCollection("col")
    }
    val singleObservable: SingleObservable[InsertOneResult] = collection.insertOne(document)
    singleObservable.subscribe(new Observer[InsertOneResult] {
      override def onNext(result: InsertOneResult): Unit = println(s"onNext: $result")
      override def onError(e: Throwable): Unit = println(s"onError: $e")
      override def onComplete(): Unit = println("onComplete")
    })
  }
}

但这种实现存在一些问题:比如isDbDown()的检查是同步的,而且可能存在竞态条件(比如检查后连接马上失效),同时可变变量的使用也不够符合Scala的函数式风格。我想知道有没有更地道的Scala方式来判断连接是否不可用,并实现可靠的连接重建逻辑?


地道的Scala解决方案

1. 核心思路

Scala的风格更偏向线程安全的不可变状态管理和异步错误驱动的逻辑,我们可以抛弃提前检查的思路,转而在操作失败时触发连接重建,同时用原子类保证多线程环境下的安全。

2. 实现代码

import com.typesafe.scalalogging.LazyLogging
import org.mongodb.scala._
import org.mongodb.scala.result.InsertOneResult
import org.mongodb.scala.connection.MongoConnectionException
import java.util.concurrent.atomic.AtomicReference

object MongoFactory extends LazyLogging {
  private val uri: String = "mongodb+srv://*********"
  // 用AtomicReference原子管理连接实例,避免多线程竞态
  private val clientRef: AtomicReference[MongoClient] = new AtomicReference(MongoClient(uri))
  private val collectionRef: AtomicReference[MongoCollection[Document]] = 
    new AtomicReference(clientRef.get().getDatabase("db").getCollection("col"))

  // 私有方法:原子性重建连接
  private def rebuildConnection(): Boolean = {
    try {
      val oldClient = clientRef.get()
      val newClient = MongoClient(uri)
      // 原子更新:只有当前实例还是oldClient时才替换,避免并发冲突
      if (clientRef.compareAndSet(oldClient, newClient)) {
        oldClient.close() // 关闭旧连接,避免资源泄漏
        val newCollection = newClient.getDatabase("db").getCollection("col")
        collectionRef.set(newCollection)
        logger.info("MongoDB连接已成功重建")
        true
      } else {
        // 其他线程已经完成了连接重建,无需重复操作
        false
      }
    } catch {
      case e: Throwable =>
        logger.error("重建MongoDB连接失败", e)
        false
    }
  }

  // 带重试机制的插入方法
  def insertDocument(document: Document, retryTimes: Int = 2): Unit = {
    val collection = collectionRef.get()
    collection.insertOne(document)
      .subscribe(new Observer[InsertOneResult] {
        override def onNext(result: InsertOneResult): Unit = 
          logger.info(s"文档插入成功: 插入ID=${result.getInsertedId}")
        
        override def onError(e: Throwable): Unit = {
          logger.error(s"文档插入失败: ${e.getMessage}", e)
          // 仅当是连接类错误且还有重试次数时,触发重建+重试
          if (retryTimes > 0 && isConnectionRelatedError(e)) {
            logger.info(s"尝试重建连接,剩余重试次数: ${retryTimes - 1}")
            if (rebuildConnection()) {
              insertDocument(document, retryTimes - 1)
            } else {
              logger.error("无法重建连接,终止重试")
            }
          } else {
            logger.error("重试次数耗尽或非连接类错误,操作终止")
          }
        }

        override def onComplete(): Unit = 
          logger.info("插入操作执行完成")
      })
  }

  // 精准判断连接类错误
  private def isConnectionRelatedError(e: Throwable): Boolean = {
    e match {
      case _: MongoConnectionException => true
      case _: java.net.SocketTimeoutException => true
      case _: java.net.ConnectException => true
      // 可以根据MongoDB驱动的异常类型扩展更多场景
      case _ => e.getMessage.contains("connection") || e.getMessage.contains("socket")
    }
  }
}

关键改进点

  • 线程安全:用AtomicReference管理连接实例,原子性更新避免多线程下的竞态问题,比单纯的var更可靠。
  • 按需重建:只有当插入操作抛出连接类错误时才触发重建,避免了"检查后连接立即失效"的竞态问题,也减少了不必要的同步检查开销。
  • 资源友好:重建连接时主动关闭旧连接,避免资源泄漏。
  • 精准错误匹配:直接匹配MongoDB驱动抛出的连接相关异常类型,比字符串匹配更可靠。
  • 可配置重试:允许设置重试次数,灵活控制容错逻辑。

额外优化建议

如果项目中使用函数式异步库(比如cats-effect或ZIO),可以进一步用IO类型封装连接管理和插入操作,让代码更符合函数式编程的"纯函数+资源安全"风格,同时简化异步错误处理逻辑。


内容的提问来源于stack exchange,提问作者blue-sky

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 04:07:37