Java 基础体系 · 第 60/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Java 并发集合:ConcurrentHashMap、CopyOnWrite、BlockingQueue 和边界
并发集合解决的不是“多个线程同时访问就绝对安全”这一笼统问题,而是几个更具体的问题:
- 一个线程写入的数据,另一个线程何时能看见?
- 多个线程同时修改同一个集合时,集合内部结构是否会损坏?
- 一次操作是单步原子的,还是由多个步骤组成?
- 读多写少、键值查找、任务传递、快照遍历等不同访问模式,应该使用哪种数据结构?
- 当生产速度超过消费速度时,系统如何表现:阻塞、丢弃、失败,还是无限积压?
Java 25 中常用的并发集合主要包括:
ConcurrentHashMap:并发键值映射;CopyOnWriteArrayList、CopyOnWriteArraySet:写时复制集合;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 必须能明确表示“当前没有映射”,否则无法区分:
- 键不存在;
- 键存在,但值是
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);
其抽象条件是:
containsKey 与 put 是两个独立动作,中间可能被其他线程插入;putIfAbsent 则把这个判断和修改交给映射本身完成。
4. computeIfAbsent 的执行过程和限制
考虑:
ConcurrentHashMap<String, UserProfile> profiles =
new ConcurrentHashMap<>();
UserProfile profile =
profiles.computeIfAbsent(userId, id -> loadProfile(id));
逻辑过程可以理解为:
- 查找
userId; - 如果已有非空值,返回已有值;
- 如果不存在,计算新值;
- 将新值与该键关联;
- 返回最终值。
计算函数不应返回 null,否则不会建立映射。计算函数还应短小、无副作用、避免阻塞,并且不能依赖对同一个映射进行递归修改。错误示例:
map.computeIfAbsent("a", key -> {
map.put("b", "another");
return "value";
});
这种写法违反了计算函数应当简单、独立的使用前提,可能导致异常、递归更新冲突或难以诊断的延迟。更严重的情况是:
map.computeIfAbsent(key, this::callRemoteService);
远程服务变慢时,其他线程访问相关键可能长时间等待;如果函数内部又等待一个依赖该映射的线程,还可能形成死锁或循环等待。
因此,若计算过程昂贵,常见做法是先设计明确的缓存状态,例如存放 CompletableFuture<V>,并单独处理失败、超时和重试,而不是让 computeIfAbsent 隐式承担完整的异步缓存协议。
5. 计数器:merge 与 LongAdder
简单计数可以写为:
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 提供 forEach、search、reduce 等批量操作。它们可以在满足条件时使用并行任务处理映射,但仍然不能把整个操作理解为冻结快照。
例如:
long total = map.reduceValuesToLong(
1L,
Integer::longValue,
0L,
Long::sum
);
第一个参数是并行阈值。阈值越小,越可能拆分任务;阈值很大时,更接近顺序执行。并行拆分是否有收益取决于映射规模、函数成本、公共线程池负载和 CPU 资源。对小映射或很轻的函数,拆分任务本身可能比计算更昂贵。
三、CopyOnWrite:读操作看到不可变快照
1. “写时复制”是什么
CopyOnWriteArrayList 的核心策略是:
- 读操作直接读取当前内部数组;
- 修改操作复制一份数组;
- 在新数组中完成添加、删除或替换;
- 一次性发布新数组;
- 旧数组继续供已经获得它的迭代器读取。
可以把一次添加抽象为:
其中 A_old 是旧数组,x 是新元素,A_new 是复制后追加元素的新数组。
因此,读线程不需要为了遍历而阻塞写线程。代价是每次结构修改都可能复制整个数组,空间复杂度和时间复杂度通常为:
其中 是当前元素数量。
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") 创建并发布了新数组,但旧迭代器仍然遍历旧数组。
这种行为有三个直接结果:
- 遍历稳定,不会因为并发结构修改而失败;
- 遍历过程不需要外部锁;
- 迭代器看不到创建后发生的修改。
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 |
正常返回 | 立即返回 false 或 null |
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
输出顺序可能不同,但满足以下约束:
- 每个
task-N先被生产,再被消费; END在所有任务之后进入队列;- 消费者收到
END后退出; - 中断发生时,线程恢复中断标志并结束当前流程,而不是吞掉中断。
这里的 END 是一种“毒丸”协议。多个消费者时,通常需要为每个消费者放入一个结束标记,或者设计明确的关闭协调机制;只放一个毒丸只能唤醒一个消费者。
3. 有界队列是系统边界,不只是容量参数
BlockingQueue<Task> queue = new ArrayBlockingQueue<>(1000);
容量 1000 表示系统最多暂存 1000 个尚未消费的任务。队列满后,系统必须选择一种策略:
put:阻塞生产者;offer:立即失败;offer带超时:等待有限时间后失败;- 自定义拒绝策略:丢弃、记录、降级或转移到其他系统。
这实际上是一个速度关系:
设生产速率为 ,消费速率为 。
- 当长期满足 ,队列通常不会持续增长;
- 当长期满足 ,有界队列最终会填满;
- 填满后,不存在“仍然无限吞吐且不丢数据、不阻塞、不扩容”的第四种魔法选项。
无界队列只是把问题从“生产者什么时候被阻塞”改成“内存什么时候耗尽”。例如:
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. 阻塞方法会响应中断
put、take、带超时的 offer 和 poll 都可能抛出 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 的职责,但直接决定队列系统是否可靠。
九、常见误解与可观察的失败表现
误解一:并发集合中的 get 和 put 能组成事务
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。如果还要实现“按键创建一次”“按旧值替换”“按键合并”,使用 putIfAbsent、compute、replace、merge 等原子方法,而不是拆成多个普通调用。
2. 是否读远多于写,并且读者需要稳定视图?
考虑 CopyOnWriteArrayList 或 CopyOnWriteArraySet。元素数量和写入频率必须可控,否则复制成本会超过锁竞争成本。
3. 是否需要在线程之间传递工作?
考虑 BlockingQueue。先决定容量和满载策略,再决定使用 put、offer 还是超时 offer。容量不是随意填写的数字,它代表系统愿意承受的积压。
4. 是否要求精确的一致性或持久化?
如果要求跨多个对象、数据库、外部服务的原子性,单个并发集合不够。应使用锁、事务、原子状态机、日志、消息确认或数据库条件更新等机制。
十一、诊断并发集合问题的方法
并发集合问题通常不是“集合抛出了错误”,而是表现为延迟、重复、遗漏或积压。可以从状态和路径入手:
ConcurrentHashMap
重点观察:
- 是否把
containsKey加put当成原子操作; compute函数是否执行远程调用或长时间阻塞;- 是否把
size当作精确全局状态; - 值对象是否在放入后被无同步地修改;
- 是否存在大量哈希冲突或热点键。
CopyOnWrite
重点观察:
- 写入频率和集合大小;
- 堆分配、GC 次数和暂停时间;
- 迭代器是否因为快照而读取旧配置;
- 元素对象是否实际可变;
- 是否错误地把它用于任务队列。
BlockingQueue
重点观察:
put阻塞时间;offer失败次数;- 队列长度和剩余容量趋势;
- 生产速率与消费速率;
- 消费处理失败后的重试数量;
- 线程是否卡在
take、put或外部依赖; - 关闭时是否存在未消费任务。
线程转储中,如果大量线程等待队列锁或条件变量,通常说明生产者—消费者速率失衡、容量过小、消费者被外部依赖阻塞,或者公平策略和锁竞争造成额外延迟。诊断时应同时看线程状态、队列指标、处理耗时和异常日志,不能只看某一次队列长度。
十二、最终边界:集合安全不等于系统正确
可以把三类集合的职责压缩为三个句子:
ConcurrentHashMap保护并发键值结构,并为某些按键操作提供原子语义;CopyOnWrite用写入复制换取稳定、无锁式的读快照;BlockingQueue用等待或失败处理生产者与消费者之间的速度差。
它们共同提供的是并发构件,不是完整业务协议。下面这些问题仍必须由应用设计决定:
- 一个任务处理失败后是否重试;
- 任务是否允许重复执行;
- 数据库提交成功但内存更新失败时如何恢复;
- 多个集合之间如何保持一致;
- 关闭时怎样保证不丢任务;
- 超过容量时是阻塞、拒绝、丢弃还是持久化;
- 统计值是要求精确,还是允许近似;
- 读到的是实时视图、弱一致视图还是固定快照。
当访问模式是“按键并发更新”时,优先从 ConcurrentHashMap 的原子方法开始;当访问模式是“少量写入、大量稳定遍历”时,考虑 CopyOnWrite;当问题是“线程之间传递工作并限制积压”时,使用有明确容量和关闭协议的 BlockingQueue。一旦需求超出单个集合的原子边界,就应显式引入更高层的同步、事务或可靠消息机制,而不是继续更换另一种集合来掩盖协议缺失。
系列导航与关联阅读
- 系列入口:Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Java Queue 与 Deque:优先队列、双端队列、阻塞语义和选型
- 下一篇:Java Collector 深入:归约、分组、下游收集器、并行和自定义
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论