【问题标题】:Missing updates with locks and ConcurrentHashMap缺少锁和 ConcurrentHashMap 的更新
【发布时间】:2019-02-13 05:53:45
【问题描述】:

我有一个场景,我必须维护一个Map,它可以由多个线程填充,每个线程都修改它们各自的List(唯一标识符/键是线程名称),并且当一个线程的列表大小超过固定的批量大小,我们必须将记录持久化到数据库中。

聚合器类

private volatile ConcurrentHashMap<String, List<T>>  instrumentMap = new ConcurrentHashMap<String, List<T>>();
private ReentrantLock lock ;

public void addAll(List<T> entityList, String threadName) {
    try {
        lock.lock();
        List<T> instrumentList = instrumentMap.get(threadName);
        if(instrumentList == null) {
            instrumentList = new ArrayList<T>(batchSize);
            instrumentMap.put(threadName, instrumentList);
        }

        if(instrumentList.size() >= batchSize -1){
            instrumentList.addAll(entityList);
            recordSaver.persist(instrumentList); 
            instrumentList.clear();
        } else {
            instrumentList.addAll(entityList);  
        }
    } finally {
        lock.unlock();
    }

}

每 2 分钟后运行一个单独的线程(使用相同的锁)以持久保存 Map 中的所有记录(以确保每 2 分钟后保存一些内容并且地图大小不会变得太大)

if(//Some condition) {
    Thread.sleep(//2 minutes);
    aggregator.getLock().lock();
    List<T> instrumentList = instrumentMap.values().stream().flatMap(x->x.stream()).collect(Collectors.toList());
    if(instrumentList.size() > 0) {
        saver.persist(instrumentList);
        instrumentMap .values().parallelStream().forEach(x -> x.clear());
        aggregator.getLock().unlock();
    }
}

这个解决方案在我们测试的几乎所有场景中都可以正常工作,除了有时我们会看到一些记录丢失了,即它们根本没有持久化,尽管它们被很好地添加到了地图中。

我的问题是:

  1. 这段代码有什么问题?
  2. ConcurrentHashMap 不是最好的解决方案吗?
  3. 与ConcurrentHashMap 一起使用的List 是否有问题?
  4. 我应该在这里使用ConcurrentHashMap 的计算方法吗(我认为不需要,因为ReentrantLock 已经在做同样的工作了)?

【问题讨论】:

  • 不确定丢失的记录,但如果对instrumentMap 的所有访问都由lock 保护,那么使用ConcurrentMap 没有任何好处。
  • @Slaw 我同意我没有写这个,也不想改变这个,直到我理解代码的问题。谢谢你的回答
  • 好吧,我无法在显示的代码中看到问题。虽然这并不意味着问题不存在,但适当的 minimal reproducible example 证明问题会有所帮助。要检查的一件事是发生了任何无人看管的访问。 recordSaver.persist 是否曾经以非阻塞方式将列表传递给另一个线程?我问是因为您传递了List 本身,而不是副本,这意味着非同步访问可能发生在某处。相比之下,您的每两分钟保存线程调用 saver.persist 并使用包含地图中所有展平值的“副本”。
  • @Slaw 您能否在答案中添加您的观察结果,以便我可以相信您是否有效:)
  • 是否需要将仪器存储在地图中?持久性是通过仪器列表完成的,“threadName”键似乎未使用。

标签: java multithreading concurrency


【解决方案1】:

@Slaw 在 cmets 中提供的答案起到了作用。我们让 instrumentList 实例以非同步方式逃逸,即访问/操作在列表上发生而没有任何同步。通过将副本传递给其他方法来修复相同的问题。

以下代码行是发生此问题的代码

recordSaver.persist(instrumentList); 仪器列表.clear();

这里我们允许 instrumentList 实例以非同步方式逃逸,即它被传递到另一个类 (recordSaver.persist) 以对其进行操作,但我们也在清除列表在下一行(在聚合器类中),所有这些都以非同步方式发生。无法在记录保护程序中预测列表状态...真是愚蠢的错误。

我们通过将 instrumentList 的克隆副本传递给 recordSaver.persist(...) 方法来解决此问题。这样 instrumentList.clear() 对 recordSaver 中可用的列表没有影响,以便进一步操作。

【讨论】:

  • 你应该详细解释答案,只是参考评论不会削减它。
  • 我一定会这样做的。
【解决方案2】:

我明白了,您在锁中使用了 ConcurrentHashMap 的 parallelStream。我不了解 Java 8+ 流支持,但快速搜索显示,

  1. ConcurrentHashMap 是一种复杂的数据结构,过去曾经存在并发错误
  2. 并行流必须遵守complex and poorly documented usage restrictions
  3. 您正在并行流中修改数据

基于这些信息(以及我的直觉驱动的并发错误检测器™),我敢打赌,删除对 parallelStream 的调用可能会提高代码的健壮性。另外,正如@Slaw 所说,如果所有instrumentMap 的使用已经被锁保护,你应该使用普通的HashMap 代替ConcurrentHashMap。

当然,由于您没有发布recordSaver 的代码,因此它也有可能存在错误(不一定是与并发相关的错误)。特别是,您应该确保从持久存储中读取记录的代码(用于检测记录丢失的代码)是安全、正确的,并且与系统的其余部分正确同步(最好使用健壮的,行业标准的 SQL 数据库)。

【讨论】:

  • 删除对 parallelStream 的调用 - 可能会增加健壮性,但与我们面临的错误无关(因为在我们使用它的时候从未观察到错误)。带有锁或 ConcurrentHashMap.compute 的 Hashmap 中的任何一个都足够了,我同意它在这里使用的方式是矫枉过正的......但这无论如何都不会影响记录。 @slaw 建议 - 我们错过的唯一漏洞,即允许引用转义到未同步的方法(recordSaver.persist)并在执行此操作后清除列表...意味着记录可能超出范围
【解决方案3】:

看起来这是对不需要的优化的尝试。在这种情况下,越少越好,越简单越好。在下面的代码中,仅使用了两个并发概念:synchronized 确保正确更新共享列表,final 确保所有线程看到相同的值。

import java.util.ArrayList;
import java.util.List;

public class Aggregator<T> implements Runnable {

    private final List<T> instruments = new ArrayList<>();

    private final RecordSaver recordSaver;
    private final int batchSize;


    public Aggregator(RecordSaver recordSaver, int batchSize) {
        super();
        this.recordSaver = recordSaver;
        this.batchSize = batchSize;
    }

    public synchronized void addAll(List<T> moreInstruments) {

        instruments.addAll(moreInstruments);
        if (instruments.size() >= batchSize) {
            storeInstruments();
        }
    }

    public synchronized void storeInstruments() {

        if (instruments.size() > 0) {
            // in case recordSaver works async
            // recordSaver.persist(new ArrayList<T>(instruments));
            // else just:
            recordSaver.persist(instruments);
            instruments.clear();
        }
    }


    @Override
    public void run() {

        while (true) {
            try { Thread.sleep(1L); } catch (Exception ignored) {
                break;
            }
            storeInstruments();
        }
    }


    class RecordSaver {
        void persist(List<?> l) {}
    }

}

【讨论】:

  • 我会试试这个
猜你喜欢
  • 2011-05-05
  • 1970-01-01
  • 2015-12-17
  • 1970-01-01
  • 2014-06-11
  • 1970-01-01
  • 1970-01-01
  • 2013-03-26
  • 1970-01-01
相关资源
最近更新 更多