hello大家,我是逆境不可逃

线程安全容器并不等于“把普通容器换个类名就万事大吉”。真正影响线上稳定性的,是复合操作是否原子、队列是否有界、满了以后怎样背压,以及容器的读写模型是否符合业务场景。

一、先建立全局认识:并发容器到底解决什么问题

普通集合大多不是线程安全的。例如,多个线程同时修改 HashMapArrayList,可能产生数据覆盖、更新丢失、读到中间状态,甚至破坏内部结构。

最直接的办法,是给所有操作加同一把锁:

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,线程直接交接

学习这些类时,不要只问“线程安全吗”,而要连续问六个问题:

  1. 单个方法安全,还是一整段业务逻辑也安全?
  2. 读取的是实时数据,还是某一时刻的快照?
  3. 队列为空时,是返回 null、抛异常,还是阻塞等待?
  4. 队列满时,是失败、阻塞,还是继续占用内存?
  5. 是否保证 FIFO?如果按优先级排序,同优先级是否稳定?
  6. 放进线程池以后,会限制任务、增加线程,还是堆积内存?

二、ConcurrentHashMap:单个操作安全,不代表组合逻辑安全

2.1 为什么不能在并发环境直接使用 HashMap

HashMapput 不只是给数组赋值,它可能涉及哈希定位、链表插入、红黑树调整和扩容迁移。多个线程无保护地同时修改,可能出现更新丢失或结构异常。

并发键值存储通常应该使用:

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

它保证 getputremoveputIfAbsentcompute 等单次调用在各自语义下是线程安全的,但不会自动把用户在方法外拼出来的一组操作变成原子操作。

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

需要原子地基于旧值计算新值时,应使用 computemerge

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);
}

sizeput 不是一个原子整体,多个线程都可能通过检查。此外,并发计数通常采用分散更新和汇总方式,适合监控与估算,不适合实现“绝不能超过 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);
}

快照迭代器不支持 removesetadd 等修改操作,否则会抛出 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;
        }
    }
}

这段源码不需要逐行背诵,抓住三点即可:

  1. tail 只是帮助快速定位尾部,不要求每一刻都精确指向最后节点。
  2. 真正决定入队成功的是把某个节点的 nextnull CAS 成新节点。
  3. 发现其他线程已经推进结构后,当前线程继续沿链表寻找,不需要锁住整条队列。

出队采用类似思想:找到第一个仍包含元素的节点,用 CAS 把其元素置空;节点和元素的逻辑删除、物理断链可能分步完成。

4.4 为什么不要频繁调用 size

链式并发队列没有一个每次入队、出队都严格维护的廉价总数。size() 通常需要遍历节点:

  • 时间复杂度为 O(n)。
  • 遍历时其他线程仍在修改,结果只是调用期间的弱一致观察。

判断是否有元素应优先使用 isEmpty(),消费任务直接使用 poll() 的返回值,不要先 size() > 0poll()

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、超时 offerputpolltake。选哪一个,本质上是在确定过载策略。

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 中断不能被悄悄吞掉

puttake 都会响应中断。捕获 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
  • 使用 notEmptynotFull 两个 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、ReentrantLockCondition 的典型落地:

队列满 → 生产者进入 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 选型对比

对比项ArrayBlockingQueueLinkedBlockingQueue
底层结构数组链表
容量构造时必须指定可指定;默认近似无界
一把锁入队、出队两把锁
节点开销较小每个元素一个节点
并行性入队出队竞争同一锁入队出队可一定程度并行
适合场景强调容量与内存可控吞吐敏感且明确设置容量

不要只根据“双锁吞吐更高”机械选择。真实结果还受任务大小、生产消费比例、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

需要注意四点:

  1. 它是逻辑无界队列,put 通常不会因为“满”而等待,仍有内存堆积风险。
  2. 同优先级元素不保证严格 FIFO;需要稳定顺序时,把递增序号加入比较规则。
  3. 高优先级任务持续进入时,低优先级任务可能长期饥饿。
  4. 元素入队后不要修改参与比较的字段,否则堆结构不会自动重新排序。

它的阻塞主要发生在消费者 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。应该使用 computemergeputIfAbsent 等复合原子方法。

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
  不存储元素,生产者消费者直接交接
  在线程池中会促使任务直接找线程

工程原则
  线程安全容器不保证任意业务组合原子
  优先明确容量、超时、中断和拒绝策略
  长期生产快于消费时,任何内存队列都救不了系统

十四、掌握程度自测

如果能独立回答以下问题,说明本周内容已经形成体系:

  1. 为什么 ConcurrentHashMap.get 后再 put 仍会丢更新?
  2. JDK 8 的 ConcurrentHashMap.putVal 遇到空桶、冲突桶和扩容节点分别怎么做?
  3. CopyOnWriteArrayList 为什么适合监听器,却不适合日志列表?
  4. ConcurrentLinkedQueue.poll()BlockingQueue.take() 的行为有什么本质差异?
  5. ArrayBlockingQueue 的两个 Condition 分别等待什么?
  6. LinkedBlockingQueue 的默认构造为什么可能带来线上风险?
  7. PriorityBlockingQueue 怎样保证同优先级 FIFO?
  8. DelayQueue 为什么不能单独承担重要订单超时关闭?
  9. SynchronousQueue 如何影响线程池扩容?
  10. 如何从容量、延迟、内存和拒绝策略设计一个稳定的异步任务队列?

结语

并发容器的价值,不只是“多个线程调用不会立刻出错”,而是为不同并发模型提供清晰的协作协议:ConcurrentHashMap 解决键值并发访问,写时复制解决读多写少,CAS 队列解决非阻塞传递,阻塞队列解决等待与背压,特殊队列则提供优先级、延迟和直接交接。

实际开发中最值得坚持的判断顺序是:先确定一致性边界,再确定读写模型,然后确定容量与过载行为,最后才是比较某个实现的理论性能。容器选对了,代码会简单很多;容量和语义选错了,再高性能的类也可能变成线上故障的放大器。
在这里插入图片描述

更多推荐