Java案例如何实现阻塞队列?

wen python案例 1

本文目录导读:

Java案例如何实现阻塞队列?

  1. 实现原理
  2. 手写实现(面试常考)
  3. 关键设计要点
  4. 使用Java内置的BlockingQueue
  5. 面试扩展:用synchronized实现

这是一个经典的Java并发编程案例,阻塞队列(BlockingQueue)是一种支持生产者-消费者模式的线程安全队列,当队列为空时,消费者线程会阻塞等待;当队列满时,生产者线程会阻塞等待。

下面我将从实现原理手写实现(基于ReentrantLock和Condition)以及使用Java内置API三个维度来讲解。


实现原理

阻塞队列的核心需要解决两个问题:

  1. 线程安全:多个线程同时读写队列时数据不混乱(需要锁保护)
  2. 等待/通知机制:队列满时生产者等待,队列空时消费者等待(需要条件变量)

Java中实现上述机制主要有两种方式:

  • synchronized + wait/notify
  • ReentrantLock + Condition(更灵活,推荐)

下面我们手写一个基于ReentrantLock + Condition的有界阻塞队列。


手写实现(面试常考)

import java.util.LinkedList;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
public class MyBlockingQueue<T> {
    // 底层存储结构
    private final LinkedList<T> queue;
    // 队列容量
    private final int capacity;
    // 锁
    private final Lock lock;
    // 队列非空条件(消费者等待)
    private final Condition notEmpty;
    // 队列非满条件(生产者等待)
    private final Condition notFull;
    public MyBlockingQueue(int capacity) {
        this.capacity = capacity;
        this.queue = new LinkedList<>();
        this.lock = new ReentrantLock();
        this.notEmpty = lock.newCondition();
        this.notFull = lock.newCondition();
    }
    /**
     * 往队列尾部添加元素,如果队列满则阻塞等待
     */
    public void put(T element) throws InterruptedException {
        lock.lock();  // 获取锁(可中断)
        try {
            // 当队列满时,生产者线程在notFull条件上等待
            while (queue.size() == capacity) {
                notFull.await();
            }
            // 生产数据
            queue.addLast(element);
            // 通知消费者:队列非空了
            notEmpty.signal();
        } finally {
            lock.unlock();
        }
    }
    /**
     * 从队列头部取出元素,如果队列空则阻塞等待
     */
    public T take() throws InterruptedException {
        lock.lock();
        try {
            // 当队列空时,消费者线程在notEmpty条件上等待
            while (queue.size() == 0) {
                notEmpty.await();
            }
            // 消费数据
            T element = queue.removeFirst();
            // 通知生产者:队列非满了
            notFull.signal();
            return element;
        } finally {
            lock.unlock();
        }
    }
    /**
     * 返回队列当前大小(调试用)
     */
    public int size() {
        lock.lock();
        try {
            return queue.size();
        } finally {
            lock.unlock();
        }
    }
}

使用示例(生产者-消费者)

public class BlockingQueueDemo {
    public static void main(String[] args) {
        // 创建一个容量为3的阻塞队列
        MyBlockingQueue<Integer> queue = new MyBlockingQueue<>(3);
        // 生产者线程
        Thread producer = new Thread(() -> {
            try {
                for (int i = 1; i <= 10; i++) {
                    queue.put(i);
                    System.out.println("生产: " + i);
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        });
        // 消费者线程
        Thread consumer = new Thread(() -> {
            try {
                for (int i = 1; i <= 10; i++) {
                    Integer value = queue.take();
                    System.out.println("消费: " + value);
                    Thread.sleep(100);  // 模拟消费耗时
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        });
        producer.start();
        consumer.start();
    }
}

关键设计要点

为什么用 while 而不是 if 进行条件判断?

// 必须使用 while 循环,不能使用 if
while (queue.size() == capacity) {
    notFull.await();
}

原因:线程被唤醒后,条件可能再次不满足(虚假唤醒或另一个生产者抢先填充了队列),使用 while 循环可以重新检查条件,这是 Guarded Suspension(保护性暂停) 模式的要求。

signal() vs signalAll()

  • signal():只唤醒一个等待线程(更高效),当前场景中,生产者只唤醒一个消费者,消费者只唤醒一个生产者,所以用 signal() 足够了。
  • signalAll():唤醒所有等待线程(安全性更高但性能较差)。

锁的选择

特性 synchronized ReentrantLock
可中断 不支持 支持
超时等待 不支持 支持
多个条件变量 不直接支持 支持(newCondition()
公平性 非公平 可设置公平/非公平

对于阻塞队列这种需要两个条件变量(notEmpty和notFull)的场景,ReentrantLock + Condition 是更自然的选择。


使用Java内置的BlockingQueue

JDK已经提供了完善的阻塞队列实现,实际开发中直接使用即可:

import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
public class BuiltInBlockingQueueDemo {
    public static void main(String[] args) throws InterruptedException {
        // 创建容量为3的有界阻塞队列
        BlockingQueue<Integer> queue = new ArrayBlockingQueue<>(3);
        // 生产者
        new Thread(() -> {
            try {
                for (int i = 1; i <= 10; i++) {
                    queue.put(i);
                    System.out.println("生产: " + i);
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }).start();
        // 消费者
        new Thread(() -> {
            try {
                for (int i = 1; i <= 10; i++) {
                    Integer value = queue.take();
                    System.out.println("消费: " + value);
                    Thread.sleep(100);
                }
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }).start();
    }
}

Java内置的阻塞队列实现(基于 ReentrantLock):

实现类 特性
ArrayBlockingQueue 有界、数组结构、FIFO
LinkedBlockingQueue 可有界/无界、链表结构
PriorityBlockingQueue 无界、支持优先级排序
SynchronousQueue 容量为0(直接传递)

面试扩展:用synchronized实现

如果面试官不让用ReentrantLock,可以用 synchronized + wait/notify

public class SyncBlockingQueue<T> {
    private final LinkedList<T> queue = new LinkedList<>();
    private final int capacity;
    public SyncBlockingQueue(int capacity) {
        this.capacity = capacity;
    }
    public synchronized void put(T element) throws InterruptedException {
        while (queue.size() == capacity) {
            wait();
        }
        queue.addLast(element);
        notifyAll();  // 必须用 notifyAll,防止只唤醒同类线程
    }
    public synchronized T take() throws InterruptedException {
        while (queue.isEmpty()) {
            wait();
        }
        T element = queue.removeFirst();
        notifyAll();
        return element;
    }
}

注意:用 synchronized 时,只有一个条件变量(wait/notify 是全局的),所以必须使用 notifyAll(),否则可能发生“信号丢失”(生产者唤醒了生产者,消费者没被唤醒)。


维度 说明
核心机制 锁 + 条件等待/通知
推荐实现 ReentrantLock + 两个Condition
关键细节 while 循环检查条件、signal() 精确唤醒
生产使用 直接使用 ArrayBlockingQueue / LinkedBlockingQueue

如果这是面试场景,建议按以下思路回答:

  1. 先讲清楚阻塞队列核心功能(满则等、空则等)
  2. 手写核心代码(put/take)
  3. 解释为什么用while不用if
  4. 对比synchronized和ReentrantLock方案的区别

抱歉,评论功能暂时关闭!