【问题标题】:ConcurrentHashMap[String, AtomicInteger] or ConcurrentHashMap[String, Int] for thread-safe counters?ConcurrentHashMap[String, AtomicInteger] 或 ConcurrentHashMap[String, Int] 用于线程安全计数器?
【发布时间】:2020-10-13 22:31:52
【问题描述】:

当通过ConcurrentHashMap 中的键递增并发计数器时,使用常规Int 作为值是否安全,还是必须使用AtomicInteger?例如考虑以下两种实现

ConcurrentHashMap[String, Int]

final class ExpensiveMetrics(implicit system: ActorSystem, ec: ExecutionContext) {
  import scala.collection.JavaConverters._
  private val chm = new ConcurrentHashMap[String, Int]().asScala

  system.scheduler.schedule(5.seconds, 60.seconds)(publishAllMetrics())

  def countRequest(key: String): Unit =
    chm.get(key) match {
      case Some(value) => chm.update(key, value + 1)
      case None => chm.update(key, 1)
    }

  private def resetCount(key: String) = chm.replace(key, 0)

  private def publishAllMetrics(): Unit =
    chm foreach { case (key, value) =>
      // publishMetric(key, value.doubleValue())
      resetCount(key)
    }
}

ConcurrentHashMap[String, AtomicInteger]

final class ExpensiveMetrics(implicit system: ActorSystem, ec: ExecutionContext) {
  import scala.collection.JavaConverters._
  private val chm = new ConcurrentHashMap[String, AtomicInteger]().asScala

  system.scheduler.schedule(5.seconds, 60.seconds)(publishAllMetrics())

  def countRequest(key: String): Unit =
    chm.getOrElseUpdate(key, new AtomicInteger(1)).incrementAndGet()
  
  private def resetCount(key: String): Unit =
    chm.getOrElseUpdate(key, new AtomicInteger(0)).set(0)

  private def publishAllMetrics(): Unit =
    chm foreach { case (key, value) =>
      // publishMetric(key, value.doubleValue())
      resetCount(key)
    }
}

以前的实现安全吗?如果没有,在 sn-p 中的什么时候可以引入竞争条件,为什么?


问题的上下文是 AWS CloudWatch 指标,如果在每个请求上发布,这些指标在高频 API 上可能会变得非常昂贵。所以我正在尝试将它们“批处理”并定期发布。

【问题讨论】:

标签: scala concurrency thread-safety counter scala-collections


【解决方案1】:

第一个实现不正确,因为countRequest 方法不是原子的。考虑一下这一系列事件:

  • 线程 A 和 B 都调用 countRequest,键为 "foo"
  • 线程A获取计数器值,我们称之为x
  • 线程 B 获取计数器值。它是同一个值 x,因为线程 A 还没有更新计数器。
  • 线程 B 使用新的计数器值 x+1 更新映射
  • 线程 A 更新映射,因为它在 B 写入新的计数器值之前获得了计数器值,所以它也写入了 x+1。

计数器应该是 x+2,但它是 x+1。这是一个经典的丢失更新问题。

由于使用了 `getOrElseUpdate` 方法,第二个实现也有类似的问题。 `ConcurrentHashMap` 没有该方法,因此 Scala 包装器需要模拟它。我认为实现是从 `scala.collection.mutable.MapOps` 继承的,它的定义如下: ``` def getOrElseUpdate(key: K, op: => V): V = 获取(键)匹配{ 案例一些(v)=> v case None => val d = op;这(键)= d; d } ``` 这显然不是原子的。

要正确实现这一点,请在ConcurrentHashMap 上使用compute 方法。

此方法将自动执行,因此您不需要AtomicInteger

【讨论】:

  • 我认为concurrent.Map#getOrElseUpdate 的实现使用putIfAbsent,文档状态为atomic
  • 你是对的。正如我所说,JConcurrentMapWrapper 确实扩展了MapOps,并从那里继承了一个非原子的getOrElseUpdate 实现。但它继承自scala.collection.concurrent.Map,它使用实际上是原子的实现来覆盖该实现。尽管如此,我认为compute 仍然更好,因为使用AtomicInteger 您需要同步两次,一次在地图上,一次在AtomicInteger 上。
猜你喜欢
  • 2012-08-20
  • 2014-03-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-09-22
  • 2011-04-15
  • 1970-01-01
  • 2016-07-16
相关资源
最近更新 更多