【问题标题】:how to avoid race condition when using Scala's Actor使用 Scala 的 Actor 时如何避免竞争条件
【发布时间】:2011-04-23 06:21:04
【问题描述】:

我正在编写一段代码,当缓冲区(列表)增长到一定大小时,它将填充 mongoDB 集合。

import scala.actors.Actor
import com.mongodb.casbah.Imports._
import scala.collection.mutable.ListBuffer

class PopulateDB extends Actor {
  val buffer = new ListBuffer[DBObject]
  val mongoConn = MongoConnection()
  val mongoCol = mongoConn("casbah_test")("logs")

  def add(info: DBObject = null) {
    if (info != null) buffer += info

    if (buffer.size > 0 && (info == null || buffer.length >= 1000)) {
      mongoCol.insert(buffer.toList)
      buffer.clear
      println("adding a batch")
    }
  }

  def act() {
    loop {
      react {
        case info: DBObject => add(info)

        case msg if msg == "closeConnection" =>
          println("Close connection")
          add()
          mongoConn.close
      }
    }
  }
}

但是,当我运行以下代码时,scala 偶尔会在“mongoCol.insert(buffer.toList)”行上抛出“ConcurrentModificationException”。我很确定它与“mongoCol.insert”有关。我想知道代码是否有任何根本性的错误。或者我应该使用 Akka 的“atomic {...}”之类的东西来避免这个问题。

这是完整的堆栈跟踪:

PopulateDB@7e859a68: caught java.util.ConcurrentModificationException
java.util.ConcurrentModificationException
    at java.util.LinkedHashMap$LinkedHashIterator.nextEntry(LinkedHashMap.java:373)
    at java.util.LinkedHashMap$EntryIterator.next(LinkedHashMap.java:392)
    at java.util.LinkedHashMap$EntryIterator.next(LinkedHashMap.java:391)
    at org.bson.BSONEncoder.putObject(BSONEncoder.java:113)
    at org.bson.BSONEncoder.putObject(BSONEncoder.java:67)
    at com.mongodb.DBApiLayer$MyCollection.insert(DBApiLayer.java:215)
    at com.mongodb.DBApiLayer$MyCollection.insert(DBApiLayer.java:180)
    at com.mongodb.DBCollection.insert(DBCollection.java:85)
    at com.mongodb.casbah.MongoCollectionBase$class.insert(MongoCollection.scala:561)
    at com.mongodb.casbah.MongoCollection.insert(MongoCollection.scala:864)
    at PopulateDB.add(PopulateDB.scala:14)
    at PopulateDB$$anonfun$act$1$$anonfun$apply$1.apply(PopulateDB.scala:26)
    at PopulateDB$$anonfun$act$1$$anonfun$apply$1.apply(PopulateDB.scala:25)
    at scala.actors.ReactorTask.run(ReactorTask.scala:34)
    at scala.actors.Reactor$class.resumeReceiver(Reactor.scala:129)
    at PopulateDB.scala$actors$ReplyReactor$$super$resumeReceiver(PopulateDB.scala:5)
    at scala.actors.ReplyReactor$class.resumeReceiver(ReplyReactor.scala:69)
    at PopulateDB.resumeReceiver(PopulateDB.scala:5)
    at scala.actors.Actor$class.searchMailbox(Actor.scala:478)
    at PopulateDB.searchMailbox(PopulateDB.scala:5)
    at scala.actors.Reactor$$anonfun$startSearch$1$$anonfun$apply$mcV$sp$1.apply(Reactor.scala:114)
    at scala.actors.Reactor$$anonfun$startSearch$1$$anonfun$apply$mcV$sp$1.apply(Reactor.scala:114)
    at scala.actors.ReactorTask.run(ReactorTask.scala:36)
    at scala.concurrent.forkjoin.ForkJoinPool$AdaptedRunnable.exec(ForkJoinPool.java:611)
    at scala.concurrent.forkjoin.ForkJoinTask.quietlyExec(ForkJoinTask.java:422)
    at scala.concurrent.forkjoin.ForkJoinWorkerThread.mainLoop(ForkJoinWorkerThread.java:340)
    at scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:325)

谢谢, 德里克

【问题讨论】:

  • 你能发布堆栈跟踪吗?什么组件抛出了并发修改异常?
  • 谢谢。我已经用堆栈跟踪更新了这个问题。
  • 'buffer' 和 'add' 公开是有原因的吗?我会让它们成为私有/受保护的,以确保它们不会在其他地方被意外调用(由非参与者线程)。
  • 我同意。它们也应该是私有的。但就我而言,我确信它们不是异常的原因
  • 只是想加入讨论:我遇到了同样的问题,这是因为我在我的存储库的构造函数中调用了 com.mongodb.casbah.commons.conversions.scala.RegisterConversionHelpers() .该函数不应该与 mongo 查询同时执行,因为它会修改 encodingHooks 哈希映射,这不是线程安全的。在所有 mongo 调用解决问题之前删除该调用或仅执行一次

标签: scala mongodb actor casbah


【解决方案1】:

DBObject 不是线程安全的;您正在发送一个 DBObject 与您的演员消息。它可能稍后会再次修改,这将导致并发修改问题。

我建议首先尝试在 DBObject 上使用clone(),因为它进入actor,然后将其放入缓冲区。它只是一个浅拷贝,但至少应该足以在支持 DBObject 上的键的 LinkedHashMap 上引起并发修改问题(通过 LHM,它保持有序)。

我会尝试:

  def act() {
    loop {
      react {
        case info: DBObject => add(info.clone())

        case msg if msg == "closeConnection" =>
          println("Close connection")
          add()
          mongoConn.close
      }
    }
  }

如果这不起作用,请在将 DBObject 发送到 Actor 后查看您正在修改的其他任何地方。

【讨论】:

  • 似乎克隆不适用于 DBObject,这令人惊讶,因为我认为克隆是从 AnyRef 继承的。我最终使用了 Synchronized{}。
  • 在 DBObject 上不可用,但在 BasicDBObject 上可用。除非您创建了自己的 DBObject 版本,否则 BasicDBObject 是 Java 驱动程序中提供的具体实现(并由 Casbah 使用)并且可以使用 clone()。
【解决方案2】:

为什么在下面class

class PopulateDB extends Actor

您是否保留多个PupulateDB 演员?我希望object PopulateDB extends Actor,这样一个演员就可以专注于这项任务。

除此之外,问题似乎出在 casbah 或 mongodb 本身。

【讨论】:

  • 问题似乎出在 Java 驱动程序内部(Casbah 封装了该驱动程序)。看起来需要同步一些东西,而不是......虽然我不了解 Scala 演员以及 Akka 的
  • @Brendan 设置 vals private 以确保它们不在其他地方使用,但似乎同时对 mongoCol.insert 的两次调用会导致问题。除了将 PopulateDB 设置为单例之外,actor 架构无能为力,以确保同时不能有多个调用。
  • @Brendan 哎呀,对不起。我以为你是德里克。
  • @DCSobral 不用担心......我也很困惑,最后一天花了一些时间寻找原因(我通常使用 Akka 作为演员)。事实证明 DBObject 根本不是线程安全的......这会导致这种情况。我已经在研究 DBObject 的不可变版本,这进一步突出了需求。
  • 感谢您的意见。将代码更改为同步 { mongoCol.insert(temp) } 解决了我的问题
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-06-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多