【问题标题】:Kotlin concurrency for ConcurrentHashMapConcurrentHashMap 的 Kotlin 并发
【发布时间】:2020-10-11 17:36:57
【问题描述】:

我正在尝试在定期清除的 hashmap 上支持并发。我有一个缓存可以存储一段时间的数据。每 5 分钟后,此缓存中的数据会发送到服务器。刷新后,我想清除缓存。问题是当我刷新时,数据可能会被写入此映射,而我正在使用现有密钥执行此操作。我将如何使这个进程线程安全?

data class A(val a: AtomicLong, val b: AtomicLong) {
   fun changeA() {
      a.incrementAndGet()
   }
}

class Flusher {
   private val cache: Map<String, A> = ConcurrentHashMap()
   private val lock = Any()
   fun retrieveA(key: String){
       synchronized(lock) {
          return cache.getOrPut(key) { A(key, 1) }
       }
   }
 
   fun flush() {
      synchronized(lock) {
           // send data to network request
           cache.clear()
      }
   }
}

// Existence of multiple classes like CacheChanger
class CacheChanger{
  fun incrementData(){
      flusher.retrieveA("x").changeA()
  }
}

我担心上面的缓存没有正确同步。有没有更好/正确的方法来锁定这个缓存,这样我就不会丢失数据?我应该创建缓存的深层副本并清除它吗?

既然上面的数据可能会被另一个changer修改,那会不会导致问题?

【问题讨论】:

  • 除了retrieve和flush之外,还有其他修改地图的功能吗?这两个函数都是在同一个锁上同步的,你怕什么问题?
  • 另外,如果你的所有访问都是同步的,为什么还要使用 ConcurrentHashMap?
  • ConcurrentHashMap 本身是线程安全的。扩展方法getOrPut 似乎也是线程安全的(基于文档)。如果没有任何其他方法以非线程安全的方式修改映射 - 您可以摆脱这个锁。
  • 问题是A类的值可以改变。如果 A 类的值发生更改并且我将其清除怎么办。我将更新示例。
  • @michalik OP 摆脱锁并不安全,因为刷新需要是原子的 - 需要读取然后清除整个映射,并且不能与此过程交错写入。

标签: java kotlin concurrency java.util.concurrent thread-synchronization


【解决方案1】:

最简单的解决方案(但性能较差)是完全依赖锁。

您可以将ConcurrentHashMap 更改为普通的HashMap

然后你必须直接在函数retrieve中应用你的所有更改:

fun retrieveA(key: String, mod: (A) -> Unit): A {
    synchronized(lock) {
        val obj: A = cache.getOrPut(key) { A(key, 1) }
        mod(obj)
        cache.put(obj)
        return obj
    }
}

我希望它能编译(我不是 Kotlin 专家)。

然后你像这样使用它:

class CacheChanger {
    fun incrementData() {
        flusher.retrieveA("x") { it.changeA() }
    }
}

好吧,我承认这段代码不是真正的 Kotlin ;) 你应该使用 Kotlin lambda 而不是 Consumer 接口。我已经有一段时间没有玩过 Kotlin 了。如果有人能解决它,我将非常感激。

【讨论】:

    【解决方案2】:

    你可以摆脱锁。

    在 flush 方法中,不是读取整个地图(例如通过迭代器)然后清除它,而是一个一个地删除每个元素。

    我不确定你是否可以使用迭代器的 remove 方法(我稍后会检查),但你可以使用 keyset 对其进行迭代,并为每个键调用 cache.remove() - 这将给出您存储的值并以原子方式将其从缓存中删除。

    棘手的部分是如何确保 A 类的对象在通过网络发送之前不会被修改...您可以这样做:

    当您通过retrieveA 获取一些x 并修改对象时,您需要确保它仍在缓存中。只需再调用一次retrieve。如果你得到完全相同的对象,那很好。如果不同,则表示该对象已被删除并通过网络发送,但您不知道是否也发送了修改,或者修改之前的对象状态已发送。不过,我认为在您的情况下,您可以简单地重复整个过程(应用更改并检查对象是否相同)。但这取决于您的应用程序的具体情况。

    如果您不想增加两次,那么在通过网络发送数据时,您必须读取计数器 a 的内容,将其存储在某个局部变量中,然后将 a 减少该数量(通常它会变为零)。然后在CacheChanger 中,当您从第二次检索中获得不同的对象时,您可以检查该值是否为零(考虑了您的修改),或者非零,这意味着您的修改只是几分之一秒迟到了,你必须重复这个过程。

    您也可以将incrementAndGet 替换为compareAndSwap,但这可能会产生稍差的性能。在这种方法中,您尝试将大于 1 的值交换,而不是递增。在通过网络发送之前,您尝试将值交换为 -1 以表示该值无效。如果第二次交换失败,则意味着有人同时更改了该值,您需要再检查一次以通过网络发送最新的值,然后循环重复该过程(仅在交换到-1 成功)。在交换到大一的情况下,您还可以循环重复该过程,直到交换成功。如果失败,则意味着其他人交换了更大的值,或者Flusher 交换了-1。在后一种情况下,您知道必须再调用一次retrieveA 才能获取新对象。

    【讨论】:

    • ConcurrentHashMap::clear 方法已经提供并保证与ConcurrentHashMap 中的其他方法一起保证线程安全时,为什么需要遍历所有条目并模仿它的行为?您基本上放弃了 concurenthashmap 的优点,转而支持外部同步。
    • @michalik 因为 clear 不返回地图的内容。您想以原子方式执行此操作,获取内容并同时删除它们。
    • @michalk 看到上面的评论
    • 对,我错过了关于通过网络发送数据的评论。在这种情况下,您的解决方案将起作用,但 OPs 解决方案(带锁)也将起作用(只要需要保留这两个操作的原子性),但在修改值时他还需要持有锁(例如通过公开Flusher 中用于修改的方法)。此外,您所描述的解决方案(似乎是一种 CAS 算法)可能需要CacheChanger 在数据发生变化时进行 CAS 检查。但总的来说,这实际上取决于应用程序的具体情况。
    • @michalk 是的,OP 的替代方案是保留锁,将映射更改为常规映射(不需要 ConcurrentHashMap),但在方法检索期间应用所有修改,同时仍然保持锁定。
    猜你喜欢
    • 1970-01-01
    • 2020-10-24
    • 2017-01-14
    • 2022-11-06
    • 1970-01-01
    • 1970-01-01
    • 2018-12-20
    • 2016-08-02
    • 1970-01-01
    相关资源
    最近更新 更多