You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

ConcurrentHashMap迭代弱一致性问题:多线程发布订阅丢失历史键

ConcurrentHashMap迭代的弱一致性问题

问题现象

  • 单个发布线程持续向ConcurrentHashMap中顺序添加新键(key-0、key-1...key-n)
  • 多个订阅线程尝试读取Map中所有现有值时,会出现早期添加的键丢失的情况:
    • 例如订阅线程读取到Map的size为100时,理论上应包含key-0到key-99,但实际无法获取全部键;
    • 甚至在key-100已经添加后,仍会随机丢失部分早期键。

最初猜测是哈希碰撞链处理导致全量读取无法返回碰撞相关键,实际是ConcurrentHashMap的弱一致性特性导致该问题。

测试代码说明

提供三个测试变体验证问题:

  1. test():稳定复现键丢失问题
  2. testWithNoCollisions():初始化足够大的Map(几乎无哈希碰撞),偶尔能获取完整键集合
  3. testWithLocks():通过外部锁阻塞写操作,可稳定获取完整键集合

原本期望ConcurrentHashMap能提供某一时刻的快照视图,但测试表明其迭代仅保证弱一致性,无法满足强一致性的快照需求。

完整测试代码

import java.util.HashSet;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Collectors;
import java.util.stream.IntStream;

public class ConcurrentHashMapWeakConsistency {

    public static void main(String[] args) throws InterruptedException {
        test();
//        testWithNoCollisions();// Sometimes it passes when there is no collision
//        testWithLocks();// This one works but it has external synchronize mechanism.
    }

    public static void test() throws InterruptedException {
        var map = new ConcurrentHashMap<String, Integer>();

        var publisher = Executors.newSingleThreadExecutor();
        var totalMessages = 1_000_000;
        publisher.submit(() -> {
            for (int i = 0; i < totalMessages; i++) {
                var key = "key-" + i;
                map.put(key, i);
                if (i % 10000 == 0) {
                    System.out.printf("Published %d messages.\n", i);
                }
            }
            System.out.printf("Published all %d messages\n", totalMessages);
        });

        var subscriberCount = 100;
        var subscribers = Executors.newFixedThreadPool(10);
        var subsLatch = new CountDownLatch(subscriberCount);

        IntStream.range(0, 100).parallel().forEach(subscriberId -> subscribers.submit(() -> {
            try {
                var existingKeys = new HashSet<>(map.keySet());
                var size = existingKeys.size();

                //Note the keys are inserted by publisher in sequential order.
                // Hence, existing keys values should have all keys from range 0 to size-1
                // This is where weak consistency shows up.
                var missingKeys = IntStream.range(0, size).filter(num -> !existingKeys.contains("key-" + num))
                        .sorted()
                        .boxed()
                        .toList();

                if (!missingKeys.isEmpty()) {
                    var sortedExistingKeys = existingKeys.stream()
                            .sorted((k1, k2) -> Integer.parseInt(k1.split("-")[1]) - Integer.parseInt(k2.split("-")[1]))
                            .toList();
                    var start = sortedExistingKeys.get(0);
                    var end = sortedExistingKeys.get(sortedExistingKeys.size() - 1);
                    var missingStart = "key-" + missingKeys.get(0);
                    var missingEnd = "key-" + missingKeys.get(missingKeys.size() - 1);
                    throw new RuntimeException(String.format("Subscriber %d missing %d keys! Result key range is [ %s, %s]. However, many numbers in range [ %s, %s ] are missing", subscriberId, missingKeys.size(), start, end, missingStart, missingEnd));
                }
            } catch (Exception ex) {
                ex.printStackTrace();
                System.exit(-1);
            } finally {
                subsLatch.countDown();
                System.out.printf("Subscriber %d finished reading from map.\n", subscriberId);
            }
        }));


        subsLatch.await();
        publisher.shutdownNow();
        subscribers.shutdownNow();
    }

    public static void testWithNoCollisions() throws InterruptedException {
        var map = new ConcurrentHashMap<String, Integer>(1_000_000, .25f);

        var publisher = Executors.newSingleThreadExecutor();
        var totalMessages = 10_000;
        publisher.submit(() -> {
            for (int i = 0; i < totalMessages; i++) {
                var key = "key-" + i;
                map.put(key, i);
                if (i % 10000 == 0) {
                    System.out.printf("Published %d messages.\n", i);
                }
            }
            System.out.printf("Published all %d messages\n", totalMessages);
        });

        var subscriberCount = 100;
        var subscribers = Executors.newFixedThreadPool(10);
        var subsLatch = new CountDownLatch(subscriberCount);

        IntStream.range(0, 100).parallel().forEach(subscriberId -> subscribers.submit(() -> {
            try {
                var existingKeys = new HashSet<>(map.keySet());
                var size = existingKeys.size();

                //Note the keys are inserted by publisher in sequential order.
                // Hence, existing keys values should have all keys from range 0 to size-1
                // This is where weak consistency shows up.
                var missingKeys = IntStream.range(0, size).filter(num -> !existingKeys.contains("key-" + num))
                        .sorted()
                        .boxed()
                        .toList();

                if (!missingKeys.isEmpty()) {
                    var sortedExistingKeys = existingKeys.stream()
                            .sorted((k1, k2) -> Integer.parseInt(k1.split("-")[1]) - Integer.parseInt(k2.split("-")[1]))
                            .toList();
                    var start = sortedExistingKeys.get(0);
                    var end = sortedExistingKeys.get(sortedExistingKeys.size() - 1);
                    var missingStart = missingKeys.get(0);
                    var missingEnd = missingKeys.get(missingKeys.size() - 1);
                    throw new RuntimeException(String.format("Subscriber %d missing %d keys! Result key range is [ %s, %s]. However, many numbers in range [ %s, %s ] are missing", subscriberId, missingKeys.size(), start, end, missingStart, missingEnd));
                }
            } catch (Exception ex) {
                ex.printStackTrace();
                System.exit(-1);
            } finally {
                subsLatch.countDown();
                System.out.printf("Subscriber %d finished reading from map.\n", subscriberId);
            }
        }));


        subsLatch.await();
        publisher.shutdownNow();
        subscribers.shutdownNow();
    }

    public static void testWithLocks() throws InterruptedException {
        var map = new ConcurrentHashMap<String, Integer>();

        var publisher = Executors.newSingleThreadExecutor();
        var totalMessages = 1_000_000;
        var subscriberActive = new AtomicBoolean(false);
        publisher.submit(() -> {
            for (int i = 0; i < totalMessages; i++) {
                var key = "key-" + i;
                // get the subscriber lock
                while (!subscriberActive.compareAndSet(false, true)) ;
                map.put(key, i);
                subscriberActive.compareAndSet(true, false);
                if (i % 10000 == 0) {
                    System.out.printf("Published %d messages.\n", i);
                }
            }
            System.out.printf("Published all %d messages\n", totalMessages);
        });

        var subscriberCount = 100;
        var subscribers = Executors.newFixedThreadPool(10);
        var subsLatch = new CountDownLatch(subscriberCount);

        IntStream.range(0, 100).parallel().forEach(subscriberId -> subscribers.submit(() -> {
            try {
                while (!subscriberActive.compareAndSet(false, true)) ;
                var existingKeys = new HashSet<>(map.keySet());
                subscriberActive.compareAndSet(true, false);

                var size = existingKeys.size();

                //Note the keys are inserted by publisher in sequential order.
                // Hence, existing keys values should have all keys from range 0 to size-1
                // This is where weak consistency shows up.
                var missingKeys = IntStream.range(0, size).filter(num -> !existingKeys.contains("key-" + num))
                        .sorted()
                        .mapToObj(i -> Integer.valueOf(i))
                        .collect(Collectors.toList());

                if (!missingKeys.isEmpty()) {
                    var sortedExistingKeys = existingKeys.stream()
                            .sorted((k1, k2) -> Integer.parseInt(k1.split("-")[1]) - Integer.parseInt(k2.split("-")[1]))
                            .toList();
                    var start = sortedExistingKeys.get(0);
                    var end = sortedExistingKeys.get(sortedExistingKeys.size() - 1);
                    var missingStart = missingKeys.get(0);
                    var missingEnd = missingKeys.get(missingKeys.size() - 1);
                    throw new RuntimeException(String.format("Subscriber %d missing %d keys! Result key range is [ %s, %s]. However, many numbers in range [ %s, %s ] are missing", subscriberId, missingKeys.size(), start, end, missingStart, missingEnd));
                }
            } catch (Exception ex) {
                ex.printStackTrace();
                System.exit(-1);
            } finally {
                subsLatch.countDown();
                System.out.printf("Subscriber %d finished reading from map.\n", subscriberId);
            }
        }));


        subsLatch.await();
        publisher.shutdownNow();
        subscribers.shutdownNow();
    }
}

原因解释

ConcurrentHashMap的迭代器是弱一致性的:

  • 迭代器创建后,会反映创建时或之后的Map状态,但不会抛出ConcurrentModificationException;
  • 在迭代过程中,若Map发生扩容或哈希链修改,迭代器可能会跳过部分元素(尤其是旧哈希桶中的元素),也可能会包含新添加的元素;
  • 当存在哈希碰撞时,扩容或链更新的过程中,旧链的部分元素可能还未被迁移到新桶,此时迭代器遍历到旧桶时可能会遗漏这些元素,这也是test()方法稳定复现问题的核心原因;
  • testWithNoCollisions()中Map初始容量足够大,避免了扩容和哈希链的频繁修改,因此偶尔能获取完整集合,但仍不保证强一致性;
  • testWithLocks()通过外部锁实现了读写互斥,强制保证了读取时的快照一致性,因此不会出现键丢失。

总结:ConcurrentHashMap的设计目标是高并发下的性能,而非强一致性快照。若需要获取某一时刻的完整快照,需额外加锁,或使用ConcurrentHashMap的copyOf方法(Java 17+)来生成不可变副本。


内容的提问来源于stack exchange,提问作者telu

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 05:37:03