判断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
相关产品推荐
相关产品推荐

