Java 并发容器一次讲透:ConcurrentHashMap、CopyOnWriteArrayList 与七种队列
hello大家,我是逆境不可逃
线程安全容器并不等于“把普通容器换个类名就万事大吉”。真正影响线上稳定性的,是复合操作是否原子、队列是否有界、满了以后怎样背压,以及容器的读写模型是否符合业务场景。
一、先建立全局认识:并发容器到底解决什么问题
普通集合大多不是线程安全的。例如,多个线程同时修改 HashMap、ArrayList,可能产生数据覆盖、更新丢失、读到中间状态,甚至破坏内部结构。
最直接的办法,是给所有操作加同一把锁:
synchronized (lock) {
map.put(key, value);
}
这样通常能保证正确性,但所有读写都挤在一把锁上,并发度较低。JUC 并发容器会根据数据结构特点,组合使用:
volatile:保证可见性
CAS:无锁竞争更新
synchronized / ReentrantLock:保护局部结构
分段或分桶:缩小竞争范围
不可变快照:让读取无需加锁
Condition:实现队列非空、非满等待
先把本文的类放进一张地图:
并发容器
│
├── Map
│ └── ConcurrentHashMap:高并发键值读写
│
├── List
│ └── CopyOnWriteArrayList:读多写极少,读取快照
│
└── Queue
├── ConcurrentLinkedQueue:无界、非阻塞、FIFO
└── BlockingQueue:支持阻塞等待
├── ArrayBlockingQueue:数组、有界、FIFO
├── LinkedBlockingQueue:链表、可选容量、FIFO
├── PriorityBlockingQueue:按优先级、逻辑无界
├── DelayQueue:到期后才能获取、逻辑无界
└── SynchronousQueue:容量为 0,线程直接交接
学习这些类时,不要只问“线程安全吗”,而要连续问六个问题:
- 单个方法安全,还是一整段业务逻辑也安全?
- 读取的是实时数据,还是某一时刻的快照?
- 队列为空时,是返回
null、抛异常,还是阻塞等待? - 队列满时,是失败、阻塞,还是继续占用内存?
- 是否保证 FIFO?如果按优先级排序,同优先级是否稳定?
- 放进线程池以后,会限制任务、增加线程,还是堆积内存?
二、ConcurrentHashMap:单个操作安全,不代表组合逻辑安全
2.1 为什么不能在并发环境直接使用 HashMap
HashMap 的 put 不只是给数组赋值,它可能涉及哈希定位、链表插入、红黑树调整和扩容迁移。多个线程无保护地同时修改,可能出现更新丢失或结构异常。
并发键值存储通常应该使用:
ConcurrentHashMap<String, Integer> map = new ConcurrentHashMap<>();
它保证 get、put、remove、putIfAbsent、compute 等单次调用在各自语义下是线程安全的,但不会自动把用户在方法外拼出来的一组操作变成原子操作。
2.2 最大陷阱:get 后 put 仍然会丢失更新
下面的代码看似用了线程安全容器:
Integer oldValue = map.get("apple");
Integer newValue = oldValue == null ? 1 : oldValue + 1;
map.put("apple", newValue);
问题在于它实际包含“读、计算、写”三步。两个线程可能同时读到 0,然后都写入 1,最终少算一次。
初始 apple = 0
线程 A:get -> 0
线程 B:get -> 0
线程 A:put -> 1
线程 B:put -> 1
预期为 2,实际为 1
需要原子地基于旧值计算新值时,应使用 compute 或 merge:
map.merge("apple", 1, Integer::sum);
map.compute("apple", (key, oldValue) ->
oldValue == null ? 1 : oldValue + 1);
2.3 Demo:多线程词频统计
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.LongAdder;
public class WordCountDemo {
public static void main(String[] args) throws InterruptedException {
List<String> words = List.of(
"java", "juc", "java", "map",
"juc", "java", "queue", "juc");
ConcurrentHashMap<String, LongAdder> counter =
new ConcurrentHashMap<>();
CountDownLatch done = new CountDownLatch(words.size());
for (String word : words) {
new Thread(() -> {
counter.computeIfAbsent(word, key -> new LongAdder())
.increment();
done.countDown();
}).start();
}
done.await();
counter.entrySet().stream()
.sorted(Map.Entry.comparingByKey())
.forEach(entry -> System.out.println(
entry.getKey() + "=" + entry.getValue().sum()));
}
}
参考输出:
java=3
juc=3
map=1
queue=1
这里有两个关键点:
computeIfAbsent原子地完成“不存在则创建计数器”。LongAdder把高竞争计数拆散到多个单元,统计时再求和,热点竞争下通常比单个AtomicLong更容易扩展。
它适合访问次数、请求数等统计指标,但 sum() 是统计语义,不应被当作转账余额之类要求瞬时强一致的值。
2.4 常用原子复合方法
| 方法 | 语义 | 常见场景 |
|---|---|---|
putIfAbsent(k, v) | 不存在才放入 | 注册、懒加载占位 |
remove(k, v) | 当前映射仍等于指定值才删除 | 条件删除 |
replace(k, old, new) | 当前值匹配才替换 | CAS 风格更新 |
compute(k, fn) | 基于旧值计算新值 | 计数、状态更新 |
computeIfAbsent(k, fn) | 不存在才计算 | 分组容器、缓存初始化 |
merge(k, v, fn) | 无旧值则放入,有旧值则合并 | 词频统计 |
注意:传给 compute 的函数应短小,避免慢查询、远程调用或复杂嵌套更新。它可能占用桶级同步区域,执行过慢会拖住访问相关桶的线程。
2.5 为什么 ConcurrentHashMap 不允许 null
HashMap.get(key) 返回 null 时,既可能代表“key 不存在”,也可能代表“key 存在但 value 为 null”。单线程里可以再调用 containsKey 判断;并发环境中,两次调用之间映射可能已经变化。
因此 ConcurrentHashMap 禁止 null key 和 null value,让 get 返回 null 明确表示“当前没有映射”,减少并发语义歧义。
2.6 JDK 7 与 JDK 8 的结构差异
JDK 7:Segment 数组 + HashEntry 数组 + 链表
通过分段锁降低整张表的竞争
JDK 8:Node 数组 + 链表 / 红黑树
空桶使用 CAS,冲突桶锁住桶头节点
扩容时多个线程可以协助迁移
面试中常说的“ConcurrentHashMap 使用分段锁”,只适合描述 JDK 7。JDK 8 已经取消 Segment 作为主结构。
2.7 put 的关键源码主线
下面是 JDK 8+ putVal 的理解版主干:
final V putVal(K key, V value, boolean onlyIfAbsent) {
if (key == null || value == null) {
throw new NullPointerException();
}
int hash = spread(key.hashCode());
for (Node<K,V>[] tab = table;;) {
Node<K,V> first;
int n, index;
if (tab == null || (n = tab.length) == 0) {
tab = initTable();
} else if ((first = tabAt(tab, index = (n - 1) & hash)) == null) {
if (casTabAt(tab, index, null,
new Node<>(hash, key, value))) {
break; // 空桶,CAS 放入成功
}
} else if (first.hash == MOVED) {
tab = helpTransfer(tab, first); // 协助扩容
} else {
synchronized (first) {
if (tabAt(tab, index) == first) {
// 遍历链表或红黑树:更新旧值或插入新节点
}
}
// 达到阈值时可能树化
break;
}
}
addCount(1L, binCount);
return oldValue;
}
主线可以记成:
计算 hash
↓
表未初始化 → 初始化
↓
桶为空 → CAS 插入
↓
遇到 ForwardingNode → 协助扩容
↓
桶非空 → 锁桶头,处理链表或红黑树
↓
更新元素数量,必要时触发扩容
为什么锁住桶头后还要检查 tabAt(tab, index) == first?因为等待锁的过程中,扩容或其他结构变化可能已经替换了该桶,必须确认自己锁住的仍是当前桶头。
2.8 get 为什么通常不加锁
理解版源码如下:
public V get(Object key) {
int hash = spread(key.hashCode());
Node<K,V>[] tab = table;
Node<K,V> first;
int n;
if (tab != null && (n = tab.length) > 0
&& (first = tabAt(tab, (n - 1) & hash)) != null) {
if (first.hash == hash && key.equals(first.key)) {
return first.value;
}
if (first.hash < 0) {
return first.find(hash, key); // 树或扩容转发节点
}
for (Node<K,V> e = first.next; e != null; e = e.next) {
if (e.hash == hash && key.equals(e.key)) {
return e.value;
}
}
}
return null;
}
节点引用和值通过 volatile 等可见性机制安全发布,读取可以沿当前结构查找,而不必和所有写线程争同一把锁。这种读取强调并发可见性,不代表把多次 get 拼起来就能获得整个 Map 的事务快照。
2.9 size 能否做严格容量控制
不能这样写:
if (map.size() < 1000) {
map.put(key, value);
}
size 与 put 不是一个原子整体,多个线程都可能通过检查。此外,并发计数通常采用分散更新和汇总方式,适合监控与估算,不适合实现“绝不能超过 1000 个元素”的强约束。
严格容量控制需要把“检查和占位”纳入同一个同步协议,例如使用锁、信号量,或把容量交给有界队列管理。
三、CopyOnWriteArrayList:写入复制,读取快照
3.1 用通讯录理解写时复制
假设系统维护一份事件监听器列表:注册和注销很少发生,但每次事件到来都会遍历所有监听器。
CopyOnWriteArrayList 的做法是:
读线程:直接读取当前数组,不加锁
写线程:加锁
→ 复制一份新数组
→ 在新数组上修改
→ 用新数组替换旧数组
→ 解锁
正在遍历旧数组的线程不受影响,后续新读取会看到新数组。这是一种以昂贵写入换取简单快速读取的设计。
3.2 Demo:遍历期间修改列表
import java.util.concurrent.CopyOnWriteArrayList;
public class CopyOnWriteDemo {
public static void main(String[] args) throws InterruptedException {
CopyOnWriteArrayList<String> list =
new CopyOnWriteArrayList<>();
list.add("A");
list.add("B");
list.add("C");
Thread reader = new Thread(() -> {
for (String value : list) {
System.out.println("遍历到:" + value);
if ("A".equals(value)) {
list.add("D");
}
}
});
reader.start();
reader.join();
System.out.println("遍历结束后的列表:" + list);
}
}
稳定输出:
遍历到:A
遍历到:B
遍历到:C
遍历结束后的列表:[A, B, C, D]
为什么本次遍历没有看到 D?因为迭代器创建时保存的是旧数组快照。添加 D 后,list 已经指向新数组,但当前迭代器仍遍历旧数组。
这也解释了它为什么不会抛出普通 ArrayList 常见的 ConcurrentModificationException:迭代器根本没有在被修改的那份数组上遍历。
3.3 关键源码
写入主线可以简化为:
public boolean add(E element) {
synchronized (lock) {
Object[] oldArray = getArray();
int length = oldArray.length;
Object[] newArray = Arrays.copyOf(oldArray, length + 1);
newArray[length] = element;
setArray(newArray);
return true;
}
}
不同版本可能使用内置监视器或显式锁,但思想没有变化:写操作串行化,复制完成后一次替换数组引用。
迭代器则保存快照:
public Iterator<E> iterator() {
return new COWIterator<>(getArray(), 0);
}
快照迭代器不支持 remove、set、add 等修改操作,否则会抛出 UnsupportedOperationException。
3.4 适合与不适合的场景
适合:
- 事件监听器、回调列表。
- 很少变化的黑白名单、规则列表。
- 配置快照、路由节点快照。
- 数据量较小,且允许读线程看到稍旧快照。
不适合:
- 日志列表、消息列表等频繁追加场景。
- 元素很多、单次复制成本高的列表。
- 内存紧张的服务。
- 要求读线程立刻看到最新写入的强实时业务。
每次写入都复制整个数组,连续写入还会制造大量短命数组,增加内存带宽和 GC 压力。它不是“线程安全版 ArrayList”的无脑替代品。
四、ConcurrentLinkedQueue:空了立即返回,不负责等待
4.1 它是什么
ConcurrentLinkedQueue 是基于链表和 CAS 实现的线程安全、无界、非阻塞 FIFO 队列。
三个关键词:
无界:没有业务容量上限,生产过快可能持续占用内存
非阻塞:方法不会为了等待非空或非满而挂起线程
FIFO:正常按入队顺序出队
这里的“非阻塞”是算法概念,表示线程竞争失败后通常重试,不依赖互斥锁让线程排队等待;它不等于每次操作都只执行一次 CAS,也不等于永远没有循环重试。
4.2 Demo:生产与批量拉取
import java.util.concurrent.ConcurrentLinkedQueue;
public class ConcurrentLinkedQueueDemo {
public static void main(String[] args) throws InterruptedException {
ConcurrentLinkedQueue<String> queue =
new ConcurrentLinkedQueue<>();
Thread producer1 = new Thread(() -> queue.offer("task-A"));
Thread producer2 = new Thread(() -> queue.offer("task-B"));
producer1.start();
producer2.start();
producer1.join();
producer2.join();
String task;
while ((task = queue.poll()) != null) {
System.out.println("处理:" + task);
}
System.out.println("再次获取:" + queue.poll());
}
}
一种可能的输出:
处理:task-A
处理:task-B
再次获取:null
两个生产线程的先后顺序可能交换,但每个已成功入队元素只会被一次 poll 移除。队列为空时,poll() 直接返回 null,消费者不会等待新任务。
4.3 offer 与 poll 的源码思想
入队的核心思想是找到合适尾节点,用 CAS 把新节点链接到尾部:
public boolean offer(E e) {
Node<E> newNode = new Node<>(Objects.requireNonNull(e));
for (Node<E> tailNode = tail, p = tailNode;;) {
Node<E> next = p.next;
if (next == null) {
if (casNext(p, null, newNode)) {
if (p != tailNode) {
casTail(tailNode, newNode);
}
return true;
}
} else if (p == next) {
p = (tailNode != (tailNode = tail)) ? tailNode : head;
} else {
p = (p != tailNode && tailNode != (tailNode = tail))
? tailNode : next;
}
}
}
这段源码不需要逐行背诵,抓住三点即可:
tail只是帮助快速定位尾部,不要求每一刻都精确指向最后节点。- 真正决定入队成功的是把某个节点的
next从nullCAS 成新节点。 - 发现其他线程已经推进结构后,当前线程继续沿链表寻找,不需要锁住整条队列。
出队采用类似思想:找到第一个仍包含元素的节点,用 CAS 把其元素置空;节点和元素的逻辑删除、物理断链可能分步完成。
4.4 为什么不要频繁调用 size
链式并发队列没有一个每次入队、出队都严格维护的廉价总数。size() 通常需要遍历节点:
- 时间复杂度为 O(n)。
- 遍历时其他线程仍在修改,结果只是调用期间的弱一致观察。
判断是否有元素应优先使用 isEmpty(),消费任务直接使用 poll() 的返回值,不要先 size() > 0 再 poll()。
4.5 什么时候换成 BlockingQueue
如果消费者写成下面这样:
while (queue.poll() == null) {
// 不断重试
}
队列长时间为空时会浪费 CPU。需要“没有任务就睡眠,有任务再唤醒”的生产者消费者模型,应使用 BlockingQueue.take()。
五、BlockingQueue:把等待、容量和背压封装进队列
5.1 四组方法必须分清
BlockingQueue 针对插入和移除提供四组语义:
| 操作 | 抛出异常 | 返回特殊值 | 一直阻塞 | 超时等待 |
|---|---|---|---|---|
| 插入 | add(e) | offer(e) | put(e) | offer(e, time, unit) |
| 移除 | remove() | poll() | take() | poll(time, unit) |
| 查看队头 | element() | peek() | 不支持 | 不支持 |
当队列满时:
add → 抛 IllegalStateException
offer → 返回 false
put → 一直等到有空间,期间可被中断
超时 offer → 最多等待指定时间,成功 true,超时 false
当队列空时:
remove → 抛 NoSuchElementException
poll → 返回 null
take → 一直等到有元素,期间可被中断
超时 poll → 最多等待指定时间,拿不到返回 null
业务代码通常更常使用 offer、超时 offer、put、poll、take。选哪一个,本质上是在确定过载策略。
5.2 Demo:有界生产者消费者
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
public class ProducerConsumerDemo {
private static final String STOP = "STOP";
public static void main(String[] args) throws InterruptedException {
BlockingQueue<String> queue = new ArrayBlockingQueue<>(2);
Thread consumer = new Thread(() -> {
try {
while (true) {
String task = queue.take();
if (STOP.equals(task)) {
System.out.println("消费者结束");
break;
}
System.out.println("消费:" + task);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
consumer.start();
queue.put("订单-1");
queue.put("订单-2");
queue.put("订单-3");
queue.put(STOP);
consumer.join();
}
}
一种可能的输出:
消费:订单-1
消费:订单-2
消费:订单-3
消费者结束
队列容量只有 2。如果生产者暂时跑得更快,第三次 put 会等待消费者腾出空间。这个等待不是缺陷,而是一种自然背压:下游处理不过来时,上游不能无限制造任务。
5.3 中断不能被悄悄吞掉
put 和 take 都会响应中断。捕获 InterruptedException 后,如果当前方法无法直接向上抛出,通常至少恢复中断标记:
try {
String task = queue.take();
handle(task);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
否则,上层线程池或关闭逻辑可能已经要求任务停止,而当前代码却把停止信号吃掉了。
六、ArrayBlockingQueue 与 LinkedBlockingQueue 怎么选
6.1 ArrayBlockingQueue:固定容量数组
ArrayBlockingQueue 在创建时必须指定容量:
BlockingQueue<Runnable> queue = new ArrayBlockingQueue<>(1000);
主要特点:
- 数组结构,创建后容量固定。
- 入队和出队共用一把
ReentrantLock。 - 使用
notEmpty、notFull两个Condition分别管理消费者和生产者。 - 默认非公平,也可以通过构造参数启用公平锁。
- 元素集中存储,内存局部性较好,不需要为每个元素创建链表节点。
put 的关键源码非常直观:
public void put(E e) throws InterruptedException {
Objects.requireNonNull(e);
final ReentrantLock lock = this.lock;
lock.lockInterruptibly();
try {
while (count == items.length) {
notFull.await();
}
enqueue(e);
} finally {
lock.unlock();
}
}
入队后会唤醒等待“非空”条件的消费者:
private void enqueue(E e) {
items[putIndex] = e;
if (++putIndex == items.length) {
putIndex = 0;
}
count++;
notEmpty.signal();
}
这正是上一周 AQS、ReentrantLock、Condition 的典型落地:
队列满 → 生产者进入 notFull 条件队列
出队后 → signal(notFull)
队列空 → 消费者进入 notEmpty 条件队列
入队后 → signal(notEmpty)
6.2 LinkedBlockingQueue:链表与双锁
BlockingQueue<Runnable> queue =
new LinkedBlockingQueue<>(1000);
主要特点:
- 链表结构,容量可以指定。
- 入队主要使用
putLock,出队主要使用takeLock。 AtomicInteger count在两把锁之间协调元素数量。- 生产和消费在一定条件下可以并行。
- 每个元素需要一个额外节点,内存开销通常高于数组队列。
关键字段的理解版如下:
private final AtomicInteger count = new AtomicInteger();
private final ReentrantLock takeLock = new ReentrantLock();
private final Condition notEmpty = takeLock.newCondition();
private final ReentrantLock putLock = new ReentrantLock();
private final Condition notFull = putLock.newCondition();
put 主线:
public void put(E e) throws InterruptedException {
Node<E> node = new Node<>(Objects.requireNonNull(e));
final ReentrantLock putLock = this.putLock;
final AtomicInteger count = this.count;
putLock.lockInterruptibly();
try {
while (count.get() == capacity) {
notFull.await();
}
enqueue(node);
int previousCount = count.getAndIncrement();
if (previousCount + 1 < capacity) {
notFull.signal();
}
} finally {
putLock.unlock();
}
// 从空变为非空,需要通知 takeLock 一侧的消费者
if (count.get() == 1) {
signalNotEmpty();
}
}
6.3 最危险的默认值:近似无界
new LinkedBlockingQueue<>() 的默认容量是 Integer.MAX_VALUE。它并非真的能放 21 亿个任务,因为通常会先耗尽堆内存。
在线程池或异步任务系统里,默认无界队列会掩盖下游过载:
请求进入速度 > 任务处理速度
↓
任务不断排队
↓
延迟越来越高
↓
任务对象占满堆内存
↓
频繁 GC,最终可能 OOM
所以工程上更稳妥的默认选择通常是“明确的有界容量 + 可观测的拒绝/降级策略”。
6.4 选型对比
| 对比项 | ArrayBlockingQueue | LinkedBlockingQueue |
|---|---|---|
| 底层结构 | 数组 | 链表 |
| 容量 | 构造时必须指定 | 可指定;默认近似无界 |
| 锁 | 一把锁 | 入队、出队两把锁 |
| 节点开销 | 较小 | 每个元素一个节点 |
| 并行性 | 入队出队竞争同一锁 | 入队出队可一定程度并行 |
| 适合场景 | 强调容量与内存可控 | 吞吐敏感且明确设置容量 |
不要只根据“双锁吞吐更高”机械选择。真实结果还受任务大小、生产消费比例、GC、缓存局部性等影响,应以压测和监控为准。
七、三种特殊阻塞队列
7.1 PriorityBlockingQueue:先处理更重要的任务
PriorityBlockingQueue 根据自然顺序或 Comparator 决定出队顺序,不保证 FIFO。
Demo:
import java.util.Comparator;
import java.util.concurrent.PriorityBlockingQueue;
public class PriorityQueueDemo {
record Task(String name, int priority) {}
public static void main(String[] args) {
PriorityBlockingQueue<Task> queue =
new PriorityBlockingQueue<>(11,
Comparator.comparingInt(Task::priority).reversed());
queue.offer(new Task("普通报表", 1));
queue.offer(new Task("支付回调", 10));
queue.offer(new Task("库存同步", 5));
while (!queue.isEmpty()) {
Task task = queue.poll();
System.out.println(task.name() + ":" + task.priority());
}
}
}
输出:
支付回调:10
库存同步:5
普通报表:1
需要注意四点:
- 它是逻辑无界队列,
put通常不会因为“满”而等待,仍有内存堆积风险。 - 同优先级元素不保证严格 FIFO;需要稳定顺序时,把递增序号加入比较规则。
- 高优先级任务持续进入时,低优先级任务可能长期饥饿。
- 元素入队后不要修改参与比较的字段,否则堆结构不会自动重新排序。
它的阻塞主要发生在消费者 take() 等待非空,不在生产者等待非满。
7.2 DelayQueue:到期后才能取出的队列
DelayQueue<E extends Delayed> 内部使用优先队列,按剩余延迟排序。即使队列中有元素,只要队头还没到期,take() 仍会等待。
Demo:
import java.util.concurrent.DelayQueue;
import java.util.concurrent.Delayed;
import java.util.concurrent.TimeUnit;
public class DelayQueueDemo {
static class DelayTask implements Delayed {
private final String name;
private final long deadlineNanos;
DelayTask(String name, long delay, TimeUnit unit) {
this.name = name;
this.deadlineNanos = System.nanoTime()
+ unit.toNanos(delay);
}
@Override
public long getDelay(TimeUnit unit) {
long remaining = deadlineNanos - System.nanoTime();
return unit.convert(remaining, TimeUnit.NANOSECONDS);
}
@Override
public int compareTo(Delayed other) {
DelayTask that = (DelayTask) other;
return Long.compare(this.deadlineNanos, that.deadlineNanos);
}
}
public static void main(String[] args) throws InterruptedException {
DelayQueue<DelayTask> queue = new DelayQueue<>();
long start = System.currentTimeMillis();
queue.put(new DelayTask("任务-2", 200, TimeUnit.MILLISECONDS));
queue.put(new DelayTask("任务-1", 100, TimeUnit.MILLISECONDS));
for (int i = 0; i < 2; i++) {
DelayTask task = queue.take();
long elapsed = System.currentTimeMillis() - start;
System.out.println(task.name + ",约 " + elapsed + " ms");
}
}
}
参考输出:
任务-1,约 101 ms
任务-2,约 201 ms
系统调度存在误差,不能期待毫秒数完全一致。
DelayQueue.take() 的源码设计还包含 leader-follower 优化:只让一个等待线程按队头剩余时间进行定时等待,其他消费者无限等待。队头变化时再通知合适线程,避免所有消费者同时定时唤醒造成惊群。
典型场景:
- 订单超时检查。
- 延迟重试。
- 本地缓存过期清理。
- 连接空闲检测。
但它只在当前 JVM 内存中。订单关闭等重要业务不能只依赖一个内存 DelayQueue:进程重启会丢任务,多实例会带来归属问题,还要考虑幂等、持久化、补偿扫描和分布式竞争。
另外应使用 System.nanoTime() 计算时间间隔,因为它是单调时间源;currentTimeMillis() 可能受系统时钟校准影响。
7.3 SynchronousQueue:容量为 0 的直接交接点
SynchronousQueue 不保存任何普通元素。每次 put 必须等到另一个线程执行匹配的 take,反过来也一样。
Demo:
import java.util.concurrent.SynchronousQueue;
public class SynchronousQueueDemo {
public static void main(String[] args) throws InterruptedException {
SynchronousQueue<String> queue = new SynchronousQueue<>();
Thread producer = new Thread(() -> {
try {
System.out.println("生产者准备交付");
queue.put("数据包");
System.out.println("生产者交付完成");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
producer.start();
Thread.sleep(200);
System.out.println("主线程开始接收");
System.out.println("收到:" + queue.take());
producer.join();
System.out.println("队列大小:" + queue.size());
}
}
参考输出:
生产者准备交付
主线程开始接收
收到:数据包
生产者交付完成
队列大小:0
put 返回时,数据已经被消费者接走,因此队列大小始终为 0。它更像接力棒交接点,而不是仓库。
无参构造默认采用非公平模式,内部偏向栈式匹配;公平模式通常使用队列式匹配:
new SynchronousQueue<>(); // 非公平
new SynchronousQueue<>(true); // 公平
Executors.newCachedThreadPool() 使用 SynchronousQueue。因为任务无法排队,若当前没有空闲线程直接接手,线程池就倾向于创建新线程,最大线程数又接近无限,因此突发流量下可能快速创建大量线程。
八、队列与线程池:选错队列会改变整个线程池
线程池执行任务的核心顺序可以简化为:
1. 工作线程少于 corePoolSize → 创建核心线程执行
2. 否则尝试把任务放入 workQueue
3. 队列放不下 → 尝试创建非核心线程,直到 maximumPoolSize
4. 仍无法接收 → 执行拒绝策略
因此队列不是一个可随意替换的配件:
| 队列 | 对线程池行为的影响 |
|---|---|
无界 LinkedBlockingQueue | 任务几乎总能入队,maximumPoolSize 通常难以生效,风险转为任务和延迟堆积 |
有界 ArrayBlockingQueue | 队列满后继续扩线程,达到上限后拒绝,容量和过载行为可控 |
SynchronousQueue | 不能缓存任务,优先直接交给线程,容易快速扩展线程数 |
PriorityBlockingQueue | 可按优先级执行,但逻辑无界,且任务需可比较或提供比较器 |
一个更可控的线程池示例:
ThreadPoolExecutor executor = new ThreadPoolExecutor(
4,
8,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(200),
runnable -> {
Thread thread = new Thread(runnable);
thread.setName("order-worker-" + thread.getId());
return thread;
},
new ThreadPoolExecutor.CallerRunsPolicy());
这个配置表达了一套清晰的过载协议:
4 个核心线程
↓ 忙不过来
最多缓存 200 个任务
↓ 队列满
最多扩到 8 个线程
↓ 仍然过载
由提交任务的线程自己执行,反向降低提交速度
CallerRunsPolicy 是一种简单背压,但不一定适合所有场景。例如提交线程是 Netty I/O 线程时,让它执行慢业务会阻塞事件循环。此时可能更适合快速失败、降级或转存到可靠消息系统。
容量也不能拍脑袋。应结合:
- 可接受的最大排队时间。
- 单个任务平均和 P99 执行时间。
- 允许占用的内存。
- 上下游超时时间。
- 峰值提交速率与处理速率。
一个粗略的延迟约束思路是:队列容量不能大到让队尾任务在开始执行前,就已经超过调用方超时时间。
九、实际开发中最常见的错误
9.1 认为“并发容器中的对象”也自动线程安全
ConcurrentHashMap<String, ArrayList<String>> groups =
new ConcurrentHashMap<>();
groups.computeIfAbsent("java", key -> new ArrayList<>())
.add("A");
Map 能安全发布和定位这个 ArrayList,但多个线程同时对同一个 ArrayList.add 仍不安全。可以根据业务改用线程安全列表、不可变替换,或对每个分组设计独立同步协议。
9.2 先判断再操作
if (!map.containsKey(key)) {
map.put(key, value);
}
改为:
map.putIfAbsent(key, value);
类似地,不要先 queue.isEmpty() 再 poll()。检查后状态随时可能变化,直接根据 poll() 返回值处理。
9.3 在 compute 回调里做慢操作
map.computeIfAbsent(userId, id -> remoteService.load(id));
远程调用可能很慢、超时或重入 Map,导致相关桶长时间受影响。缓存加载需要进一步考虑重复加载、超时、异常、缓存击穿和异步刷新,不应只靠一个复杂 compute 回调解决。
9.4 无界队列当缓冲区
无界只是“API 不规定合理上限”,不是资源无限。所有生产速度长期超过消费速度的无界队列,最终都把流量问题转换成内存和延迟问题。
9.5 用 CopyOnWriteArrayList 承接频繁写入
写一次复制一次,列表越大越昂贵。批量更新时可以先构造完整不可变快照,再一次替换引用,而不是对 COW 列表连续调用数千次 add。
9.6 修改 PriorityBlockingQueue 中的排序字段
元素进入堆后修改优先级,队列不会主动重新建堆。正确方式通常是移除旧元素再按新优先级重新入队,或让排序字段不可变。
9.7 把 DelayQueue 当可靠定时任务中心
单 JVM 内存队列无法天然处理宕机恢复、多实例协调和任务审计。核心业务应组合数据库、消息队列、调度系统或补偿任务。
9.8 吞掉 InterruptedException
中断是任务取消和服务关闭的重要协作机制。不要只打印异常后继续死循环,应向上抛出或恢复中断标记,并退出当前流程。
十、关键源码阅读路线
不建议从头逐行阅读几千行源码,可以沿调用主线看:
10.1 ConcurrentHashMap
put
→ putVal
→ initTable / casTabAt / synchronized(first)
→ treeifyBin
→ addCount
→ transfer / helpTransfer
get
→ spread
→ tabAt
→ Node.find / TreeBin.find / ForwardingNode.find
重点问题:空桶为何 CAS、冲突桶锁谁、如何识别扩容、为何可以协助迁移。
10.2 CopyOnWriteArrayList
add
→ 获取写锁
→ getArray
→ Arrays.copyOf
→ setArray
iterator
→ 获取当前 array
→ COWIterator 保存 snapshot
重点问题:写为什么贵、读为什么无需锁、迭代为何看不到后续修改。
10.3 ConcurrentLinkedQueue
offer
→ 找尾部候选节点
→ casNext 链接新节点
→ 尝试推进 tail
poll
→ 找第一个有效 item
→ CAS 将 item 置空
→ 必要时推进 head
重点问题:为什么 head/tail 可以暂时滞后、线性化点是哪次 CAS、逻辑删除与物理断链为何分离。
10.4 ArrayBlockingQueue
put → lockInterruptibly → while(满) notFull.await → enqueue
take → lockInterruptibly → while(空) notEmpty.await → dequeue
重点问题:为什么用 while、两个 Condition 各自等待什么、signal 后为什么还要重新竞争锁。
10.5 LinkedBlockingQueue
put → putLock + notFull
take → takeLock + notEmpty
count → 协调两把锁和状态转换通知
重点问题:双锁如何提升并行性、从空到非空/从满到非满时为何需要跨锁通知。
10.6 特殊队列
PriorityBlockingQueue → 二叉堆 + 扩容 + notEmpty
DelayQueue → PriorityQueue + leader + available Condition
SynchronousQueue → TransferStack / TransferQueue 直接匹配
十一、高频面试题与参考答案
11.1 ConcurrentHashMap 为什么线程安全
JDK 8+ 使用数组、链表和红黑树组织数据;空桶插入主要使用 CAS,桶内冲突更新使用桶头级别的 synchronized,配合 volatile、CAS 和安全发布保证可见性与原子更新,扩容时还能让多个线程协助迁移。它没有给整张表统一加锁。
11.2 JDK 7 与 JDK 8 的 ConcurrentHashMap 有什么区别
JDK 7 主结构是 Segment + HashEntry,通过分段锁并发;JDK 8 主结构改成 Node 数组,桶为空时 CAS,冲突桶锁桶头,链表过长时可以树化,并支持多线程协助扩容。
11.3 ConcurrentHashMap 的 get 为什么通常不加锁
节点和值通过可见性机制安全发布,读取可沿数组和节点结构查找。它提供单次读取的并发语义,但不保证多次读取组成整张 Map 的事务快照。
11.4 get 后 put 为什么不安全
因为它是两个独立的线程安全方法,中间还夹着计算,整体不是原子操作。其他线程可以在这几步之间修改同一个 key。应该使用 compute、merge、putIfAbsent 等复合原子方法。
11.5 ConcurrentHashMap 为什么不允许 null
并发环境需要让 get 返回 null 明确代表当前不存在映射。若允许 null value,就难以区分不存在和值为 null,而且通过额外检查也会遇到状态在两次调用之间变化的问题。
11.6 ConcurrentHashMap 的 size 是精确的吗
它会尽力统计当前元素数,但并发修改一直发生时,任何返回值都可能立即过时。适合监控和一般查询,不适合拿来实现严格容量检查。size() < limit 后再 put 不是原子操作。
11.7 CopyOnWriteArrayList 为什么读不加锁
写线程不修改读线程正在使用的旧数组,而是在锁内复制并修改新数组,最后替换数组引用。读线程只读取某个稳定数组快照,因此不需要与写线程争用同一把锁。
11.8 CopyOnWriteArrayList 有什么缺点
每次写入需要复制数组,耗时和内存开销随数据量增加;迭代器读取的是快照,无法保证看到最新修改;频繁写入还会制造大量临时数组和 GC 压力。
11.9 ConcurrentLinkedQueue 和 BlockingQueue 的区别
前者是非阻塞无界队列,空队列 poll 返回 null,不会等待;后者提供 put/take 等阻塞方法,可用于生产者消费者协作。是否有界要看具体 BlockingQueue 实现。
11.10 ArrayBlockingQueue 和 LinkedBlockingQueue 怎么选
数组队列固定容量、节点开销小、入队出队共用一把锁;链表队列应显式设置容量,节点开销较大,但入队和出队使用两把锁,可一定程度并行。实际选择要结合内存控制、吞吐和压测,线上通常都应明确有界。
11.11 put、offer、add 有什么区别
队列满时,put 可中断地一直等待空间,offer 立即返回 false,超时版 offer 最多等待指定时间,add 失败则抛异常。它们代表不同的过载处理策略。
11.12 PriorityBlockingQueue 是有界的吗
它是逻辑无界的,构造参数中的初始容量不是最大容量。生产过快仍可能导致内存耗尽,它的阻塞主要发生在消费者等待非空时。
11.13 PriorityBlockingQueue 能保证同优先级 FIFO 吗
不能天然保证。若业务需要稳定顺序,应在元素中加入全局递增序号,并在优先级相同时按序号比较。
11.14 DelayQueue 为什么队列有元素时 take 还可能阻塞
因为只有延迟时间小于等于 0 的元素才可被取出。若队头尚未到期,消费者需要等待剩余延迟,或等待一个截止时间更早的新元素到来。
11.15 SynchronousQueue 为什么说容量为 0
它不保存普通元素,一次插入必须与一次移除直接匹配。生产者成功返回时,消费者已经接走数据,所以 size() 始终为 0。
11.16 为什么线程池不建议随便使用无界队列
当提交速度持续高于处理速度时,任务会无限堆积,排队延迟不断上升,并占用大量内存,最终可能频繁 GC 或 OOM。无界队列还会让最大线程数通常难以生效。
11.17 什么是背压
背压是下游处理能力不足时,将压力反馈给上游,让上游减速、等待、拒绝或降级。有限容量队列加阻塞/拒绝策略,是进程内实现背压的常见方式。
11.18 阻塞队列底层为什么使用 while 而不是 if 等待
线程从 await 或类似等待中返回,只代表它有机会重新检查条件,不代表条件一定仍满足。它可能经历虚假唤醒,也可能醒来后被其他线程抢先改变状态,所以必须循环检查“非空”或“非满”条件。
11.19 并发容器的迭代器都是强一致的吗
不是。许多并发容器提供弱一致迭代:遍历不会因并发修改而失败,但可能看到部分修改。CopyOnWriteArrayList 更明确地遍历创建迭代器时的数组快照。都不能默认当成数据库事务快照。
11.20 如何选择本文中的容器
先按数据模型选 Map、List 或 Queue,再判断读写比例、是否需要等待、是否需要容量、是否要优先级或延迟语义。最后根据内存、延迟和过载策略做压测,不能只凭“线程安全”四个字选型。
十二、进阶补充:影响线上稳定性的几个细节
12.1 弱一致性不等于数据错乱
并发容器为了允许读写并行,遍历、计数等聚合观察可能不是某个瞬间的全局快照。这是明确的并发语义,不等于内部结构损坏。
如果业务必须基于完全一致的多个值做决策,就需要额外锁、不可变快照、版本号或事务机制。
12.2 单方法原子性与业务原子性是两层问题
accountMap.computeIfPresent(id, (key, account) -> {
account.setBalance(account.getBalance() - 100);
return account;
});
即使映射计算受保护,Account 对象还可能被其他路径直接访问,余额扣减也可能涉及数据库和流水记录。容器只能保证容器协议,无法自动保证整个领域事务。
12.3 队列容量本质上是延迟预算与内存预算
容量太小,会频繁拒绝突发任务;容量太大,会让任务在队列里等待到失去业务价值。容量设置要结合处理速率、允许突发量、超时阈值和单任务内存,而不是简单使用一个“业内推荐数字”。
12.4 监控不能只有 queue.size
建议同时观察:
- 任务提交速率、完成速率和拒绝速率。
- 队列当前长度、容量使用率和增长趋势。
- 排队等待时间、执行时间、端到端耗时的 P95/P99。
- 活跃线程数、线程池最大线程数。
- GC 暂停、堆内存和任务对象大小。
队列长度只是结果。真正需要回答的是:为什么堆积、还要多久耗尽容量、任务是否已经超时。
12.5 批量处理要限制单批数量
许多队列提供 drainTo,可以减少逐个获取的锁竞争:
List<Task> batch = new ArrayList<>(100);
queue.drainTo(batch, 100);
但一次抽干整个队列可能让其他消费者饥饿,也可能形成超大批次。因此应设置合理的单批上限,并正确处理批任务部分失败。
十三、一页式复盘
ConcurrentHashMap
单方法并发安全
复合更新用 compute / merge / putIfAbsent
JDK 8+:空桶 CAS,冲突桶锁桶头,支持协助扩容
不允许 null,size 不适合严格容量控制
CopyOnWriteArrayList
写时复制,读数组快照
适合读多写极少、数据量不大
写入昂贵,迭代看不到后续更新
ConcurrentLinkedQueue
无界、非阻塞、FIFO
poll 空队列返回 null
size 通常需要遍历,不适合频繁调用
ArrayBlockingQueue
数组、固定有界、一把锁、两个 Condition
强调容量和内存可控
LinkedBlockingQueue
链表、入队出队双锁
一定要警惕默认 Integer.MAX_VALUE 容量
PriorityBlockingQueue
按优先级,不保证同优先级 FIFO
逻辑无界,可能内存堆积和低优先级饥饿
DelayQueue
元素到期后才能 take
适合 JVM 内延迟任务,不等于可靠分布式调度
SynchronousQueue
不存储元素,生产者消费者直接交接
在线程池中会促使任务直接找线程
工程原则
线程安全容器不保证任意业务组合原子
优先明确容量、超时、中断和拒绝策略
长期生产快于消费时,任何内存队列都救不了系统
十四、掌握程度自测
如果能独立回答以下问题,说明本周内容已经形成体系:
- 为什么
ConcurrentHashMap.get后再put仍会丢更新? - JDK 8 的
ConcurrentHashMap.putVal遇到空桶、冲突桶和扩容节点分别怎么做? CopyOnWriteArrayList为什么适合监听器,却不适合日志列表?ConcurrentLinkedQueue.poll()和BlockingQueue.take()的行为有什么本质差异?ArrayBlockingQueue的两个Condition分别等待什么?LinkedBlockingQueue的默认构造为什么可能带来线上风险?PriorityBlockingQueue怎样保证同优先级 FIFO?DelayQueue为什么不能单独承担重要订单超时关闭?SynchronousQueue如何影响线程池扩容?- 如何从容量、延迟、内存和拒绝策略设计一个稳定的异步任务队列?
结语
并发容器的价值,不只是“多个线程调用不会立刻出错”,而是为不同并发模型提供清晰的协作协议:ConcurrentHashMap 解决键值并发访问,写时复制解决读多写少,CAS 队列解决非阻塞传递,阻塞队列解决等待与背压,特殊队列则提供优先级、延迟和直接交接。
实际开发中最值得坚持的判断顺序是:先确定一致性边界,再确定读写模型,然后确定容量与过载行为,最后才是比较某个实现的理论性能。容器选对了,代码会简单很多;容量和语义选错了,再高性能的类也可能变成线上故障的放大器。

更多推荐



所有评论(0)