Java 基础体系 · 第 60/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java 并发集合:ConcurrentHashMap、CopyOnWrite、BlockingQueue 和边界

并发集合解决的不是“多个线程同时访问就绝对安全”这一笼统问题,而是几个更具体的问题:

  1. 一个线程写入的数据,另一个线程何时能看见?
  2. 多个线程同时修改同一个集合时,集合内部结构是否会损坏?
  3. 一次操作是单步原子的,还是由多个步骤组成?
  4. 读多写少、键值查找、任务传递、快照遍历等不同访问模式,应该使用哪种数据结构?
  5. 当生产速度超过消费速度时,系统如何表现:阻塞、丢弃、失败,还是无限积压?

Java 25 中常用的并发集合主要包括:

  • ConcurrentHashMap:并发键值映射;
  • CopyOnWriteArrayListCopyOnWriteArraySet:写时复制集合;
  • BlockingQueue 及其实现:带阻塞语义的生产者—消费者队列。

它们都不能替代锁、事务、限流或业务不变量。理解它们的边界,必须先区分可见性、原子性、线程安全和业务正确性


一、并发集合到底保证什么

1. 可见性、原子性和有序性不是同一个概念

设线程 A 执行:

state.value = 42;

线程 B 随后执行:

System.out.println(state.value);

如果两者之间没有恰当的同步关系,B 不一定能按照程序员期望观察到 A 的写入。这里涉及的是可见性

如果两个线程同时执行:

count++;

这并不是一个不可分割的动作,而近似于:

int old = count;
int next = old + 1;
count = next;

两个线程可能都读到 old == 0,最后都写入 1。这里失败的是复合操作的原子性

而即使每个变量都能被看见、每个单步操作都不会破坏集合内部结构,下面这个业务条件仍可能失败:

if (!map.containsKey(id)) {
    map.put(id, value);
}

两个线程可能同时通过 containsKey,然后都执行 put。集合本身没有损坏,但“只允许首次插入”的业务不变量没有被原子地实现。

因此应区分:

概念 关注的问题
可见性 一个线程的写入能否被另一个线程观察到
原子性 一个操作是否不可分割
线程安全 并发调用是否不会破坏对象内部状态
业务原子性 多个步骤组合起来是否满足业务条件
有序性 操作或元素是否按照特定顺序发生

并发集合通常保证自己的内部状态安全,并为特定操作提供内存可见性和原子性;它不会自动把任意业务流程变成原子事务。

2. 发布对象与并发集合之间的关系

如果线程 A 把一个对象放入并发集合,随后线程 B 从该集合中成功取得这个对象,Java 并发集合的规范会为这种传递建立相应的 happens-before 关系:A 在放入对象之前对该对象的操作,可以被 B 在取得对象之后观察到。

这解决的是“对象如何安全交给另一个线程”的问题。例如:

Map<String, String> map = new ConcurrentHashMap<>();

String value = new String("ready");
// 对 value 的初始化完成
map.put("status", value);

// 另一个线程通过 map.get("status") 取得 value

但它不意味着对象在放入后就永远不可变,也不意味着对象内部所有后续字段访问都自动安全:

class Task {
    int progress; // 普通字段
}

如果一个 Task 被放进 ConcurrentHashMap 后,多个线程还会并发修改 progress,仍然需要 volatile、锁、原子变量或其他同步手段。并发集合保证的是集合操作和跨线程发布路径,不是任意元素对象的内部并发协议。


二、ConcurrentHashMap:并发键值映射

1. 它解决的核心问题

ConcurrentHashMap<K, V> 允许多个线程并发读取和更新键值对,同时避免用一把外部锁把整个映射完全串行化。

典型场景是:

ConcurrentHashMap<String, Long> counts = new ConcurrentHashMap<>();

多个线程根据键查找、插入或更新计数,而其他线程同时读取结果。

与传统 HashMap 相比:

  • HashMap 在并发结构修改时可能造成数据丢失或内部结构异常;
  • Collections.synchronizedMap(...) 通常要求所有访问通过同一个包装对象,并且遍历仍需显式同步;
  • ConcurrentHashMap 为常见查找和更新提供了更细粒度的并发控制。

Java 25 的常见实现使用哈希桶、节点和按需同步等机制;具体锁竞争细节属于实现层面,不应把某一版本的内部字段或锁布局当成 API 契约。可以依赖的是它的并发语义,而不是某个实现细节。

2. ConcurrentHashMap 不允许 null

以下代码会抛出 NullPointerException

ConcurrentHashMap<String, String> map = new ConcurrentHashMap<>();

map.put(null, "value");
map.put("key", null);

原因不仅是风格选择。对于并发映射,get(key) == null 必须能明确表示“当前没有映射”,否则无法区分:

  1. 键不存在;
  2. 键存在,但值是 null

这种明确性使得并发读取的结果更容易解释。若业务需要表示“已计算但结果为空”,应使用特殊对象、Optional(注意额外对象开销),或单独的状态类型,而不是把 null 塞进 ConcurrentHashMap

3. 单步操作与复合操作

这些方法分别是具有明确并发语义的单步操作:

map.putIfAbsent(key, value);
map.remove(key, value);
map.replace(key, oldValue, newValue);
map.computeIfAbsent(key, k -> createValue(k));
map.merge(key, 1L, Long::sum);

相反,下面的写法不是原子的:

if (!map.containsKey(key)) {
    map.put(key, value);
}

正确的“缺失时插入”应写成:

map.putIfAbsent(key, value);

其抽象条件是:

insert(k,v)={成功,若 k 不存在不改变,若 k 已存在\text{insert}(k, v) = \begin{cases} \text{成功,若 } k \text{ 不存在}\\ \text{不改变,若 } k \text{ 已存在} \end{cases}

containsKeyput 是两个独立动作,中间可能被其他线程插入;putIfAbsent 则把这个判断和修改交给映射本身完成。

4. computeIfAbsent 的执行过程和限制

考虑:

ConcurrentHashMap<String, UserProfile> profiles =
        new ConcurrentHashMap<>();

UserProfile profile =
        profiles.computeIfAbsent(userId, id -> loadProfile(id));

逻辑过程可以理解为:

  1. 查找 userId
  2. 如果已有非空值,返回已有值;
  3. 如果不存在,计算新值;
  4. 将新值与该键关联;
  5. 返回最终值。

计算函数不应返回 null,否则不会建立映射。计算函数还应短小、无副作用、避免阻塞,并且不能依赖对同一个映射进行递归修改。错误示例:

map.computeIfAbsent("a", key -> {
    map.put("b", "another");
    return "value";
});

这种写法违反了计算函数应当简单、独立的使用前提,可能导致异常、递归更新冲突或难以诊断的延迟。更严重的情况是:

map.computeIfAbsent(key, this::callRemoteService);

远程服务变慢时,其他线程访问相关键可能长时间等待;如果函数内部又等待一个依赖该映射的线程,还可能形成死锁或循环等待。

因此,若计算过程昂贵,常见做法是先设计明确的缓存状态,例如存放 CompletableFuture<V>,并单独处理失败、超时和重试,而不是让 computeIfAbsent 隐式承担完整的异步缓存协议。

5. 计数器:mergeLongAdder

简单计数可以写为:

ConcurrentHashMap<String, Long> counts = new ConcurrentHashMap<>();

counts.merge("java", 1L, Long::sum);

这里 merge 对同一个键的合并操作具有原子语义,避免了:

Long old = counts.get("java");
counts.put("java", old == null ? 1L : old + 1L);

后者会丢失并发更新。

高竞争下,常见的统计方式是:

ConcurrentHashMap<String, LongAdder> counts = new ConcurrentHashMap<>();

counts.computeIfAbsent("java", ignored -> new LongAdder())
      .increment();

long total = counts.get("java").sum();

LongAdder 通过分散竞争提高更新吞吐,但代价是:

  • sum() 不是把所有并发更新冻结后得到的线性一致快照;
  • 在其他线程持续更新时,读取值可能处于某个时间窗口;
  • 它适合统计吞吐量、访问次数等近似实时指标,不适合余额、库存、配额等必须精确判断的业务状态。

如果业务要求“扣减后绝不能低于零”,LongAdder 不是合适工具;需要原子条件更新、锁、数据库条件更新或其他事务性方案。

6. 遍历是弱一致的,不是快照

ConcurrentHashMap 的迭代器通常是弱一致的:

  • 不会因为并发修改而抛出 ConcurrentModificationException
  • 可以与更新并发进行;
  • 不保证遍历期间看到一个固定时刻的完整快照;
  • 可能看到部分更新,也可能看不到某些刚发生的更新;
  • 不提供全局遍历顺序。

例如:

for (Map.Entry<String, Integer> entry : map.entrySet()) {
    process(entry.getKey(), entry.getValue());
}

这适合“尽力处理当前可见条目”的场景,不适合生成严格一致的账单、全量导出或需要稳定分页的结果。

size()isEmpty()mappingCount() 在并发更新期间也不应被解释为全局事务快照。它们可以用于监控、近似判断或诊断,但不能单独支撑“队列是否完全为空”“是否达到精确上限”等复杂业务结论。

7. 批量操作的边界

ConcurrentHashMap 提供 forEachsearchreduce 等批量操作。它们可以在满足条件时使用并行任务处理映射,但仍然不能把整个操作理解为冻结快照。

例如:

long total = map.reduceValuesToLong(
        1L,
        Integer::longValue,
        0L,
        Long::sum
);

第一个参数是并行阈值。阈值越小,越可能拆分任务;阈值很大时,更接近顺序执行。并行拆分是否有收益取决于映射规模、函数成本、公共线程池负载和 CPU 资源。对小映射或很轻的函数,拆分任务本身可能比计算更昂贵。


三、CopyOnWrite:读操作看到不可变快照

1. “写时复制”是什么

CopyOnWriteArrayList 的核心策略是:

  • 读操作直接读取当前内部数组;
  • 修改操作复制一份数组;
  • 在新数组中完成添加、删除或替换;
  • 一次性发布新数组;
  • 旧数组继续供已经获得它的迭代器读取。

可以把一次添加抽象为:

Anew=Aold+ ⁣ ⁣+[x]A_{new} = A_{old} \mathbin{+\!\!+} [x]

其中 A_old 是旧数组,x 是新元素,A_new 是复制后追加元素的新数组。

因此,读线程不需要为了遍历而阻塞写线程。代价是每次结构修改都可能复制整个数组,空间复杂度和时间复杂度通常为:

Twrite=O(n),Stemporary=O(n)T_{\text{write}} = O(n), \qquad S_{\text{temporary}} = O(n)

其中 nn 是当前元素数量。

2. 快照迭代器的实际行为

CopyOnWriteArrayList<String> listeners =
        new CopyOnWriteArrayList<>();

listeners.add("A");

Iterator<String> iterator = listeners.iterator();

listeners.add("B");

while (iterator.hasNext()) {
    System.out.println(iterator.next());
}

输出只包含:

A

因为迭代器在创建时保存了当时的数组引用。之后的 add("B") 创建并发布了新数组,但旧迭代器仍然遍历旧数组。

这种行为有三个直接结果:

  1. 遍历稳定,不会因为并发结构修改而失败;
  2. 遍历过程不需要外部锁;
  3. 迭代器看不到创建后发生的修改。

CopyOnWriteArrayList 的迭代器不支持 Iterator.remove() 等修改操作;因为迭代器代表的是一个快照,不能把修改写回它所持有的旧数组。

3. 适合读多写少,而不是“任何并发列表”

事件监听器、路由规则、配置快照等场景通常具有如下模式:

初始化或偶尔更新:少量写入
每个请求或事件处理:大量遍历

此时复制成本由少量写操作承担,而大量读操作无需锁。

反例是高频追加队列:

CopyOnWriteArrayList<Task> tasks = new CopyOnWriteArrayList<>();

如果每秒持续追加大量任务,列表会不断复制已有数组。它既没有阻塞队列的等待语义,也没有高效的尾部生产能力,最终表现为 CPU、内存分配和垃圾回收压力。

4. 元素可变性不会自动变成快照

写时复制只保护数组结构,不复制元素对象:

class Config {
    int timeout;
}

CopyOnWriteArrayList<Config> list = new CopyOnWriteArrayList<>();
Config config = new Config();
list.add(config);

之后若某线程执行:

config.timeout = 5000;

另一个线程从列表快照中取得的仍然是同一个 Config 对象。它看到什么,取决于 timeout 字段本身的同步方式。

因此:

  • 列表结构是快照;
  • 元素引用是快照;
  • 元素对象内部状态不一定是快照。

若需要真正不可变的配置快照,应创建不可变元素,或在更新时复制元素对象,而不仅仅依赖 CopyOnWriteArrayList

5. CopyOnWriteArraySet

CopyOnWriteArraySet 基于写时复制策略提供集合语义,适合元素数量不大、读操作远多于写操作、并且需要去重的场景。

它的去重判断需要扫描已有元素,因此写入成本仍然随元素数量增长。对于大型集合或频繁变更集合,不能因为它是并发集合就忽略其线性复制和查找成本。


四、BlockingQueue:把数据传递和等待协议放进队列

1. BlockingQueue 的核心语义

BlockingQueue<E> 同时描述了容器和线程协作协议:

  • 队列有空间时,生产者可以插入;
  • 队列为空时,消费者可以等待;
  • 队列满时,生产者可以等待;
  • 等待可以被中断;
  • 某些方法还提供超时、立即失败或返回特殊值的行为。

四组方法的差异如下:

操作 成功时 失败或无条件件时
add / remove / element 正常返回 抛出异常
offer / poll / peek 正常返回 立即返回 falsenull
put / take 正常返回 必要时一直等待
offer(e, timeout) / poll(timeout) 正常返回 等待到超时后返回失败结果

null 不能作为 BlockingQueue 元素,因为 poll() 返回 null 表示“当前没有元素”。如果业务需要表达空值,应使用包装对象或显式状态。

2. 生产者—消费者的状态变化

以容量为 2 的队列为例:

初始:      []
生产 A:    [A]
生产 B:    [A, B]   —— 已满
生产 C:    put(C) 阻塞
消费 A:    [B]      —— 释放一个空间
生产 C:    [B, C]
消费 B:    [C]
消费 C:    []
再次消费:  take() 阻塞

这里的阻塞不是异常,而是背压的一种实现:消费者变慢时,生产者会被迫减速。

可以用下面的程序观察完整路径:

import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.TimeUnit;

public class BlockingQueueDemo {
    private static final String END = "<END>";

    public static void main(String[] args) throws InterruptedException {
        BlockingQueue<String> queue = new ArrayBlockingQueue<>(2);

        Thread producer = Thread.ofPlatform().name("producer").start(() -> {
            try {
                for (int i = 1; i <= 5; i++) {
                    String task = "task-" + i;
                    queue.put(task);
                    System.out.println(Thread.currentThread().getName()
                            + " produced " + task);
                }

                queue.put(END);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                System.err.println("producer interrupted");
            }
        });

        Thread consumer = Thread.ofPlatform().name("consumer").start(() -> {
            try {
                while (true) {
                    String task = queue.take();

                    if (END.equals(task)) {
                        break;
                    }

                    System.out.println(Thread.currentThread().getName()
                            + " consumed " + task);
                    TimeUnit.MILLISECONDS.sleep(100);
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                System.err.println("consumer interrupted");
            }
        });

        producer.join();
        consumer.join();
    }
}

程序需要 JDK 25 编译运行:

javac BlockingQueueDemo.java
java BlockingQueueDemo

输出顺序可能不同,但满足以下约束:

  1. 每个 task-N 先被生产,再被消费;
  2. END 在所有任务之后进入队列;
  3. 消费者收到 END 后退出;
  4. 中断发生时,线程恢复中断标志并结束当前流程,而不是吞掉中断。

这里的 END 是一种“毒丸”协议。多个消费者时,通常需要为每个消费者放入一个结束标记,或者设计明确的关闭协调机制;只放一个毒丸只能唤醒一个消费者。

3. 有界队列是系统边界,不只是容量参数

BlockingQueue<Task> queue = new ArrayBlockingQueue<>(1000);

容量 1000 表示系统最多暂存 1000 个尚未消费的任务。队列满后,系统必须选择一种策略:

  • put:阻塞生产者;
  • offer:立即失败;
  • offer 带超时:等待有限时间后失败;
  • 自定义拒绝策略:丢弃、记录、降级或转移到其他系统。

这实际上是一个速度关系:

设生产速率为 λp\lambda_p,消费速率为 λc\lambda_c

  • 当长期满足 λp<λc\lambda_p < \lambda_c,队列通常不会持续增长;
  • 当长期满足 λp>λc\lambda_p > \lambda_c,有界队列最终会填满;
  • 填满后,不存在“仍然无限吞吐且不丢数据、不阻塞、不扩容”的第四种魔法选项。

无界队列只是把问题从“生产者什么时候被阻塞”改成“内存什么时候耗尽”。例如:

BlockingQueue<Task> queue = new LinkedBlockingQueue<>();

在没有显式指定容量时,它的容量上限非常大,工程上常被当作无界队列使用。它不会自动解决生产速度大于消费速度的问题。

4. 常见实现的选择

ArrayBlockingQueue

new ArrayBlockingQueue<>(capacity)
  • 有界;
  • 基于数组;
  • 适合需要明确内存边界的生产者—消费者模型;
  • 可以选择公平策略,但公平通常会牺牲部分吞吐。

LinkedBlockingQueue

new LinkedBlockingQueue<>(capacity)
  • 可以有界,也可以不显式指定容量;
  • 基于链式节点;
  • 有界形式适合需要动态链式存储但仍要限制积压的场景;
  • 无界形式必须谨慎,因为任务积压可能转化为内存压力。

SynchronousQueue

new SynchronousQueue<>();

它没有实际容量。一次 put 必须直接与一次 take 配对,适合直接移交任务的场景:

生产者提交任务 —— 等待 —— 消费者接收任务

它不是“容量为 1 的队列”。容量为 1 的队列可以在没有消费者时暂存一个元素,而 SynchronousQueue 不能暂存。

PriorityBlockingQueue

它按优先级取出元素,而不是按插入顺序取出。它通常是无界的,元素必须具备可比较关系,或构造时提供比较器。

即使优先级相同,也不应假设它保持严格 FIFO。若业务要求同优先级按提交顺序处理,需要把序号纳入排序键。

DelayQueue

元素只有在延迟到期后才能被取出,适合延迟任务和过期对象。它不是普通 FIFO 队列,也不能直接拿来表示“提交顺序”。


五、三类集合的关键差异

类型 主要问题 读行为 写行为 是否阻塞 典型边界
ConcurrentHashMap 并发键值查找和更新 弱一致、非快照 细粒度并发更新 通常不因容量等待 复合业务操作仍需设计
CopyOnWriteArrayList 读多写少的稳定遍历 快照迭代 复制整个数组 不阻塞等待 写操作昂贵,元素仍可变
BlockingQueue 线程间传递和背压 取出会改变队列 插入会改变队列 可阻塞、可超时 容量和关闭协议必须明确

一个实际系统可能同时使用三者:

请求线程
   │
   ├── ConcurrentHashMap:保存按 ID 索引的运行状态
   │
   ├── CopyOnWriteArrayList:保存很少变更、经常遍历的监听器
   │
   └── BlockingQueue:传递待处理任务并限制积压

它们解决的是不同维度的问题,不能互相替换:

  • ConcurrentHashMap 不能自然表达“满了就等待”;
  • CopyOnWriteArrayList 不能高效表达任务队列;
  • BlockingQueue 不能直接替代按键查找;
  • 用任何一种集合都不能自动保证数据库和内存状态的一致性。

六、复合操作:线程安全集合最容易被误用的边界

1. “先检查,再操作”通常不是原子操作

错误示例:

if (queue.remainingCapacity() > 0) {
    queue.add(task);
}

检查和插入之间,其他线程可能已经填满队列。正确方式是直接使用具有明确失败语义的操作:

if (!queue.offer(task)) {
    reject(task);
}

如果需要等待一段时间:

if (!queue.offer(task, 500, TimeUnit.MILLISECONDS)) {
    rejectAfterTimeout(task);
}

同理,下面的映射代码也不安全:

if (!cache.containsKey(key)) {
    cache.put(key, value);
}

应改为:

cache.putIfAbsent(key, value);

2. “取出后修改”不是整体原子

即使队列本身是线程安全的,下面的业务过程仍可能失败:

Task task = queue.poll();
if (task != null) {
    database.save(task);
    metrics.increment();
}

队列只保证 poll 的并发安全。数据库保存失败时,任务已经从队列移除;是否重试、回滚、转移到死信队列,都需要业务协议。

一个典型失败路径是:

take 成功
  │
  ├── 处理成功 —— 确认完成
  │
  └── 处理失败 —— 任务已离开队列,必须显式重试或记录

如果希望“处理成功后再确认”,BlockingQueue 本身没有消息中间件式的 ack 语义。需要在应用层维护状态,或使用具备确认、持久化和重投能力的外部消息系统。

3. 复合条件需要外部同步或原子状态机

例如库存不能低于零:

if (stock.get() >= requested) {
    stock.addAndGet(-requested);
}

这仍然可能被并发线程插入。应使用 CAS 循环:

for (;;) {
    int current = stock.get();

    if (current < requested) {
        throw new IllegalStateException("insufficient stock");
    }

    if (stock.compareAndSet(current, current - requested)) {
        break;
    }
}

或者使用锁、数据库条件更新:

UPDATE inventory
SET quantity = quantity - ?
WHERE item_id = ?
  AND quantity >= ?;

并发集合能提供原子构件,但“多个对象之间的条件关系”通常需要更高层的同步机制。


七、内存一致性与元素对象的安全使用

1. 安全发布不等于持续安全修改

下面的初始化通常可以通过并发集合安全发布:

final class Config {
    private final String endpoint;
    private final int timeoutMillis;

    Config(String endpoint, int timeoutMillis) {
        this.endpoint = endpoint;
        this.timeoutMillis = timeoutMillis;
    }
}

final 字段和通过并发集合发布共同有助于读取初始化后的配置。

但如果对象设计为可变:

final class MutableConfig {
    int timeoutMillis;
}

放入集合后再修改普通字段,仍存在数据竞争。可以选择:

  • 让对象真正不可变,每次修改创建新对象;
  • 使用 volatile 字段;
  • 使用锁;
  • 使用原子变量;
  • 把更新操作集中到一个受控线程。

2. 不要把集合线程安全误读成“遍历期间业务稳定”

以下代码不会因为并发修改而抛出结构并发修改异常:

for (String key : map.keySet()) {
    if (shouldRemove(key)) {
        map.remove(key);
    }
}

但它不表示遍历的是固定集合,也不表示所有满足条件的键都会在这次循环中被处理。若要求精确的一致性清理,应使用明确的同步边界、批次标记、版本号,或先复制出业务需要的快照:

List<String> snapshot = List.copyOf(map.keySet());

这次复制得到的是某个时刻的近似视图,但复制完成后,snapshot 本身不再随映射变化。复制过程仍不是跨整个业务流程的数据库式快照。


八、异常、中断和关闭协议

1. 阻塞方法会响应中断

puttake、带超时的 offerpoll 都可能抛出 InterruptedException。捕获后通常应恢复中断标志:

catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    return;
}

直接吞掉异常:

catch (InterruptedException ignored) {
}

会使上层无法知道线程已经收到取消信号,导致线程继续运行、关闭变慢或任务无法停止。

2. BlockingQueue 没有内建的关闭状态

队列接口没有统一的 close()。因此必须自己定义结束协议,例如:

  • 毒丸对象;
  • 单独的停止标志;
  • 中断消费者;
  • 任务状态对象;
  • 使用更高层的执行器或消息系统。

毒丸对象必须和正常任务类型清晰区分。若使用字符串 "<END>" 作为标记,而正常业务也可能产生同样字符串,就会发生协议冲突。更稳妥的是使用专用类型:

sealed interface Work permits Task, Stop {}

record Task(String id) implements Work {}

enum Stop implements Work {
    INSTANCE
}

这样结束信号不会与普通任务值混淆。

3. 失败重试必须防止无限循环

如果消费者处理失败后立即把任务重新放回队列:

try {
    process(task);
} catch (Exception e) {
    queue.offer(task);
}

可能造成“失败—重试—失败”的忙循环。至少需要明确:

  • 最大重试次数;
  • 重试延迟;
  • 可重试异常与不可重试异常;
  • 死信存储;
  • 任务幂等性;
  • 关闭期间是否继续重试。

这些都超出了 BlockingQueue 的职责,但直接决定队列系统是否可靠。


九、常见误解与可观察的失败表现

误解一:并发集合中的 getput 能组成事务

if (map.get(key) == null) {
    map.put(key, create());
}

失败表现是重复创建、后写覆盖先写或不必要的外部调用。使用 computeIfAbsent 或显式锁定相关业务键,才能表达“缺失时创建”的原子条件。

误解二:CopyOnWriteArrayList 的读完全没有成本

读通常不需要锁,但仍然要:

  • 遍历数组;
  • 保持旧数组被迭代器引用;
  • 在写入时复制整个数组;
  • 为新数组分配内存并触发旧数组回收。

如果写入频繁,GC 压力和复制成本会成为主要问题。

误解三:有界队列能自动提升吞吐

有界队列主要提供的是积压上限和背压,不是吞吐优化。队列满时,如果生产者阻塞,调用方延迟会上升;如果选择丢弃或失败,业务数据可能丢失。必须结合业务可接受的失败方式选择。

误解四:ConcurrentHashMap.size() 可以决定全局业务动作

例如:

if (map.size() < 1000) {
    map.put(id, value);
}

多个线程都可能同时观察到小于 1000,最后超过限制。精确上限需要原子计数、锁、信号量,或把“检查和插入”设计成单一原子协议。

误解五:队列为空说明系统没有任务

在并发系统中:

if (queue.isEmpty()) {
    shutdown();
}

检查后,生产者可能立即放入任务;或者任务已经被另一个消费者取走但尚未完成。队列为空只说明某一时刻没有待取元素,不说明所有任务已经完成,也不说明生产者已经停止。


十、如何根据访问模式选择

可以先问四个问题:

1. 是否需要按键快速查找?

需要时考虑 ConcurrentHashMap。如果还要实现“按键创建一次”“按旧值替换”“按键合并”,使用 putIfAbsentcomputereplacemerge 等原子方法,而不是拆成多个普通调用。

2. 是否读远多于写,并且读者需要稳定视图?

考虑 CopyOnWriteArrayListCopyOnWriteArraySet。元素数量和写入频率必须可控,否则复制成本会超过锁竞争成本。

3. 是否需要在线程之间传递工作?

考虑 BlockingQueue。先决定容量和满载策略,再决定使用 putoffer 还是超时 offer。容量不是随意填写的数字,它代表系统愿意承受的积压。

4. 是否要求精确的一致性或持久化?

如果要求跨多个对象、数据库、外部服务的原子性,单个并发集合不够。应使用锁、事务、原子状态机、日志、消息确认或数据库条件更新等机制。


十一、诊断并发集合问题的方法

并发集合问题通常不是“集合抛出了错误”,而是表现为延迟、重复、遗漏或积压。可以从状态和路径入手:

ConcurrentHashMap

重点观察:

  • 是否把 containsKeyput 当成原子操作;
  • compute 函数是否执行远程调用或长时间阻塞;
  • 是否把 size 当作精确全局状态;
  • 值对象是否在放入后被无同步地修改;
  • 是否存在大量哈希冲突或热点键。

CopyOnWrite

重点观察:

  • 写入频率和集合大小;
  • 堆分配、GC 次数和暂停时间;
  • 迭代器是否因为快照而读取旧配置;
  • 元素对象是否实际可变;
  • 是否错误地把它用于任务队列。

BlockingQueue

重点观察:

  • put 阻塞时间;
  • offer 失败次数;
  • 队列长度和剩余容量趋势;
  • 生产速率与消费速率;
  • 消费处理失败后的重试数量;
  • 线程是否卡在 takeput 或外部依赖;
  • 关闭时是否存在未消费任务。

线程转储中,如果大量线程等待队列锁或条件变量,通常说明生产者—消费者速率失衡、容量过小、消费者被外部依赖阻塞,或者公平策略和锁竞争造成额外延迟。诊断时应同时看线程状态、队列指标、处理耗时和异常日志,不能只看某一次队列长度。


十二、最终边界:集合安全不等于系统正确

可以把三类集合的职责压缩为三个句子:

  • ConcurrentHashMap 保护并发键值结构,并为某些按键操作提供原子语义;
  • CopyOnWrite 用写入复制换取稳定、无锁式的读快照;
  • BlockingQueue 用等待或失败处理生产者与消费者之间的速度差。

它们共同提供的是并发构件,不是完整业务协议。下面这些问题仍必须由应用设计决定:

  • 一个任务处理失败后是否重试;
  • 任务是否允许重复执行;
  • 数据库提交成功但内存更新失败时如何恢复;
  • 多个集合之间如何保持一致;
  • 关闭时怎样保证不丢任务;
  • 超过容量时是阻塞、拒绝、丢弃还是持久化;
  • 统计值是要求精确,还是允许近似;
  • 读到的是实时视图、弱一致视图还是固定快照。

当访问模式是“按键并发更新”时,优先从 ConcurrentHashMap 的原子方法开始;当访问模式是“少量写入、大量稳定遍历”时,考虑 CopyOnWrite;当问题是“线程之间传递工作并限制积压”时,使用有明确容量和关闭协议的 BlockingQueue。一旦需求超出单个集合的原子边界,就应显式引入更高层的同步、事务或可靠消息机制,而不是继续更换另一种集合来掩盖协议缺失。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。