Administrator
发布于 2019-09-12 / 3527 阅读
92

BlockingQueue 家族与生产者消费者模型落地

读 ThreadPoolExecutor 源码时,我顺手压测了三种队列

9 月份啃线程池源码,看到构造函数里那个 BlockingQueue<Runnable> workQueue 参数,我才意识到自己从来没认真选过它——一直是 Executors.newFixedThreadPool(20) 一路用到底。

于是写了个压测,把常用的三种队列拉出来跑了一遍,结果挺意外。这篇记一下差异和选型。

先理清 7 组方法

BlockingQueue 继承了 Queue,所以它有两组语义不同的操作,一共 7 个方法。我第一次看的时候觉得重复,后来发现每一组对应一种"满了/空了怎么办"的策略:

抛异常返回特殊值一直阻塞超时退出
插入add(e)offer(e)put(e)offer(e, time, unit)
移除remove()poll()take()poll(time, unit)
查看element()peek()

生产代码里我用得最多的是带超时的那组。put 会一直阻塞,如果消费者挂了,生产者线程就永远卡在那,连中断都不响应(实际上响应中断,会抛 InterruptedException,但如果你 catch 了没处理就白搭)。带超时的版本至少能记个日志、走降级。

// 队列满了等 3 秒,还进不去就降级
if (!queue.offer(task, 3, TimeUnit.SECONDS)) {
    log.warn("queue full, size={}, fallback to local file", queue.size());
    dumpToFile(task);
}

ArrayBlockingQueue:一把锁管两端

数组实现,构造时就要指定容量,之后不能改。它的核心是一把锁

public class ArrayBlockingQueue<E> extends AbstractQueue<E> {
    final Object[] items;
    int takeIndex;          // 下一个要取的位置
    int putIndex;           // 下一个要放的位置
    int count;

    final ReentrantLock lock;              // 只有一把
    private final Condition notEmpty;      // 空了,消费者在这等
    private final Condition notFull;       // 满了,生产者在这等
}

put 的实现:

public void put(E e) throws InterruptedException {
    checkNotNull(e);
    final ReentrantLock lock = this.lock;
    lock.lockInterruptibly();              // 可中断地加锁
    try {
        while (count == items.length)      // 注意是 while 不是 if
            notFull.await();
        enqueue(e);
    } finally {
        lock.unlock();
    }
}

private void enqueue(E x) {
    final Object[] items = this.items;
    items[putIndex] = x;
    if (++putIndex == items.length)        // 到末尾了绕回开头,环形数组
        putIndex = 0;
    count++;
    notEmpty.signal();                     // 唤醒一个消费者
}

两个细节值得说:

  • while 而不是 ifawait() 被唤醒后必须重新检查条件,因为可能有别的线程抢先了。这是条件等待的标准写法,忽略这条会出偶发 bug。
  • items[putIndex] = x 直接写数组,没有创建对象。所以它在反复入队出队时不产生垃圾,GC 压力比链表实现小。我们做实时风控的时候特意选了它,因为它 predictable。

缺点是单锁。生产和消费抢同一把锁,高并发下会互相阻塞。

LinkedBlockingQueue:两把锁,吞吐更高

链表实现,每个元素包一个 Node。它用两把锁把入队和出队分开了:

public class LinkedBlockingQueue<E> extends AbstractQueue<E> {
    private final int capacity;
    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();
}
public void put(E e) throws InterruptedException {
    if (e == null) throw new NullPointerException();
    int c = -1;
    Node<E> node = new Node<E>(e);
    final ReentrantLock putLock = this.putLock;
    final AtomicInteger count = this.count;
    putLock.lockInterruptibly();
    try {
        while (count.get() == capacity)
            notFull.await();
        enqueue(node);
        c = count.getAndIncrement();                  // 注意是先 get 后 increment
        if (c + 1 < capacity)
            notFull.signal();                         // 还没满,再叫一个生产者进来
    } finally {
        putLock.unlock();
    }
    if (c == 0)
        signalNotEmpty();                             // 刚从空变成非空,唤醒消费者
}

为什么 c == 0 时才唤醒消费者?因为只要队列里还有元素,消费端自己会在取完一个之后继续唤醒下一个(级联通知)。只在"从空到非空"这个临界点通知一次就够了,减少无谓的 signal。这是 Doug Lea 一贯的优化思路。

入队只拿 putLock,出队只拿 takeLock,互不干扰,所以理论上吞吐是 ArrayBlockingQueue 的两倍。代价是每次入队要 new 一个 Node 对象,产生垃圾。

最大的坑在默认容量

public LinkedBlockingQueue() {
    this(Integer.MAX_VALUE);      // 21 亿
}

不传容量的话,它就是个"事实上无界"的队列。Executors.newFixedThreadPool 用的正是这个无参构造,所以任务堆积时不会触发拒绝策略(队列永远不满,maximumPoolSize 也永远用不上),一直堆到 OOM。我现在看到 newFixedThreadPool 就条件反射地想提醒一句。

SynchronousQueue:不存东西的队列

这个队列最反直觉:它的容量是 0,isEmpty() 永远返回 true,size() 永远返回 0。它不存储任何元素,只做一手交钱一手交货

public void put(E e) throws InterruptedException {
    if (e == null) throw new NullPointerException();
    if (transferer.transfer(e, false, 0) == null) {
        Thread.interrupted();
        throw new InterruptedException();
    }
}

生产者调用 put 时会阻塞,直到有一个消费者来 take;反之亦然。它内部用的是 Transferer,JDK 8 默认是非公平模式(栈实现,后进先出),构造时传 true 可以用公平模式(队列实现,先进先出)。

它的用处是把任务直接交给线程,不做缓冲Executors.newCachedThreadPool() 就是它:

public static ExecutorService newCachedThreadPool() {
    return new ThreadPoolExecutor(0, Integer.MAX_VALUE,     // 最大线程数 21 亿
                                  60L, TimeUnit.SECONDS,
                                  new SynchronousQueue<Runnable>());
}

提交任务时,如果有空闲线程在等就直接交给它;没有就新建一个线程。线程空闲 60 秒后回收。这个组合适合大量短任务、负载波动大的场景,池子会自动伸缩。

但它也意味着:没有缓冲,来多少任务就要多少线程。我见过有人用它跑批量任务,一次性提交 5000 个,直接创建 5000 个线程,机器 load 飙到 80,然后 OOM。用它必须配 maximumPoolSize

new ThreadPoolExecutor(
        8, 64, 60L, TimeUnit.SECONDS,
        new SynchronousQueue<>(),
        new ThreadPoolExecutor.CallerRunsPolicy());   // 超了就反压

压测数据

场景:4 个生产者 + 4 个消费者,每个生产者投递 100 万个 Runnable(内容是空的,只测队列本身的开销),统计总耗时。队列容量都设 10000。机器是 4 核 8 G 的 CentOS 7,JDK 8u212。

队列总耗时吞吐GC 次数YGC 总耗时
ArrayBlockingQueue(10000)8.4 s476 K ops/s00 ms
LinkedBlockingQueue(10000)5.1 s784 K ops/s37412 ms
SynchronousQueue11.7 s342 K ops/s218 ms

LinkedBlockingQueue 比 ArrayBlockingQueue 快 39%,印证了双锁的优势,代价是 37 次 YGC。SynchronousQueue 最慢,因为它每次投递都要做一次线程间的握手(park/unpark),但它不占内存、无延迟。

把容量调大到 100 万再测,两者差距缩小到 12%——因为队列越长,锁竞争的概率越低,双锁优势就不明显了。

在线程池里怎么选

队列配合的线程池适用场景
ArrayBlockingQueue固定大小 + 有界队列需要严格控制内存、任务量可预估
LinkedBlockingQueue固定大小 + 显式容量吞吐优先,任务量波动大
SynchronousQueue可伸缩 + 有最大线程数短任务、低延迟、不希望积压
PriorityBlockingQueue需要按优先级执行任务有轻重缓急,无界
DelayQueue定时/延迟任务订单超时关闭、重试队列

还有个必须知道的机制:队列类型和 maximumPoolSize 的关系。线程池的任务处理顺序是"核心线程 → 队列 → 最大线程 → 拒绝策略"。也就是说,只有队列满了之后,才会创建超过 corePoolSize 的线程。

所以如果你用无界的 LinkedBlockingQueue,队列永远不满,maximumPoolSize 这个参数就完全失效了,池子永远只有 corePoolSize 个线程。这个坑我见过两次。JDK 里 newFixedThreadPool 就是这么设计的(core = max),所以没暴露问题。

我们项目现在的标准写法:

// 订单处理:任务量可预估,控制内存优先
new ThreadPoolExecutor(
        16, 16, 0L, TimeUnit.MILLISECONDS,
        new ArrayBlockingQueue<>(5000),
        new ThreadFactoryBuilder().setNameFormat("order-pool-%d").build(),
        new ThreadPoolExecutor.CallerRunsPolicy());

// 消息推送:短任务,不希望积压,允许伸缩
new ThreadPoolExecutor(
        8, 64, 60L, TimeUnit.SECONDS,
        new SynchronousQueue<>(),
        new ThreadFactoryBuilder().setNameFormat("push-pool-%d").build(),
        new ThreadPoolExecutor.AbortPolicy());

小结

  • 四组 API 对应四种"满了怎么办"的策略:add 抛异常、offer 返回 false、put 死等、offer(超时) 等一会儿。生产环境优先用带超时的。
  • ArrayBlockingQueue 单锁 + 环形数组,不产生垃圾,GC 友好;LinkedBlockingQueue 双锁(putLock/takeLock),吞吐高 39%,但每次入队 new 一个 Node。
  • LinkedBlockingQueue 默认容量是 Integer.MAX_VALUEnewFixedThreadPool 用的就是它。任务堆积时不会触发拒绝策略,直接 OOM。
  • SynchronousQueue 容量为 0,做的是直接移交。newCachedThreadPool 用它,所以最大线程数是 21 亿,必须手动限制。
  • 线程池"核心线程 → 队列 → 最大线程"的顺序意味着:无界队列会让 maximumPoolSize 彻底失效。

看完源码最大的收获是:JDK 里这些"看起来差不多"的类,差别都藏在锁的粒度和数据结构上。选之前先想清楚自己要的是吞吐、延迟还是内存可控。

参考