ArrayBlockingQueue

mac2026-10-04  1

文章目录

简介主要属性入队列出队列总结

简介

ArrayBlockingQueue是采用数组实现的有界阻塞线程安全队列 。 线程安全是指,ArrayBlockingQueue内部通过“互斥锁”保护竞争资源,实现了多线程对竞争资源的互斥访问。

主要属性

/** 使用数组存储队列的元素 */ final Object[] items; /** 每次获取元素时的数组下标,默认为0*/ int takeIndex; /** 每次插入元素时的数组下标,默认为0 */ int putIndex; /** 队列中的数量 */ int count; /** 使用锁来保护新增元素、获取元素的入口。 */ final ReentrantLock lock; /** 队列不为空的条件 */ private final Condition notEmpty; /** 队列不满的条件 */ private final Condition notFull; // 构造函数,需要传入队列大小,默认使用非公平锁 public ArrayBlockingQueue(int capacity) { this(capacity, false); } public ArrayBlockingQueue(int capacity, boolean fair) { if (capacity <= 0) throw new IllegalArgumentException(); // 初始化数组的大小 this.items = new Object[capacity]; lock = new ReentrantLock(fair); notEmpty = lock.newCondition(); notFull = lock.newCondition(); }

入队列

// 入队列,若队列已满,则会阻塞等待 public void put(E e) throws InterruptedException { // 检查插入元素是否为空 checkNotNull(e); final ReentrantLock lock = this.lock; // 加锁,如果线程被中断,则抛出异常 lock.lockInterruptibly(); try { // 若 count == items.length ,说明队列已经满了。 // 若队列已满,则使用notFull.await()进行等待。当有线程取走元素,会调用notFull.signal()唤醒等待的线程。 // 这里之所以使用while而不是if // 是因为有可能多个线程阻塞在lock上 // 即使唤醒了可能其它线程先一步修改了队列又变成满的了 // 这时候需要再次等待 while (count == items.length) notFull.await(); // 释放锁,阻塞线程,等待该线程被唤醒 enqueue(e);// 入队 } finally { // 释放锁 lock.unlock(); } } // 元素入队 private void enqueue(E x) { final Object[] items = this.items; // 将元素存储在数组下标为putIndex位置 items[putIndex] = x; // 将putIndex加1,代表下一次插入元素的数组下标 // 如果下次元素插入数组位置等于数组的长度,说明数组已经满了,那么下一次插入的位置应该为0 if (++putIndex == items.length) putIndex = 0; count++; // 队列元素数量加1 // 插入了一个元素,队列肯定不为空了,需要通知因为获取元素而阻塞的线程。 notEmpty.signal(); } // 入队列。若队列已满,则会返回false,入队成功则会返回true public boolean offer(E e) { // 检查插入元素是否为空 checkNotNull(e); final ReentrantLock lock = this.lock; // 加锁 lock.lock(); try { // 若队列已满,则返回false if (count == items.length) return false; else { // 元素入队 enqueue(e); return true; } } finally { // 释放锁 lock.unlock(); } } // 入队列 public boolean add(E e) { // 调用父类的方法 return super.add(e); } // AbstractQueue.add(e) public boolean add(E e) { // offer 由子类ArrayBlockingQueue实现,即上面的offer方法 if (offer(e)) return true; // 入队成功则返回true else throw new IllegalStateException("Queue full"); // 入队失败则抛出异常 }

出队列

// 出队列。如果队列无元素,则阻塞等待。 public E take() throws InterruptedException { final ReentrantLock lock = this.lock; // 加锁。如果当前线程被中断,将会抛出异常。 lock.lockInterruptibly(); try { // 如果队列为空,使用notEmpty.await()进行等待。当有元素出入队列时,会调用notEmpty.signal()来唤醒等待线程。 while (count == 0) notEmpty.await(); // 释放锁,当前线程加入等待队列。 // 出队列 return dequeue(); } finally { // 释放锁 lock.unlock(); } } // 出队列 private E dequeue() { final Object[] items = this.items; @SuppressWarnings("unchecked") // 获取元素 E x = (E) items[takeIndex]; // 将数组下标为takeIndex的元素设置为空 items[takeIndex] = null; // 将takeIndex加1,代表下次获取元素的数组下标 // 如果takeIndex等于数组的长度,需要设置为0 if (++takeIndex == items.length) takeIndex = 0; // 数组元素减1 count--; if (itrs != null) itrs.elementDequeued(); // 唤醒notFull条件。元素出队列,队列肯定不满了,需要唤醒因为入队列而阻塞的线程。 notFull.signal(); return x; }

总结

ArrayBlockingQueue利用takeIndex和putIndex循环利用数组。

队列长度固定并且必须在初始化时指定。

如果出队列速度跟不上入队速度,则会导致入队线程一直阻塞。

只使用了一个锁来控制入队操作与出队操作,效率较低。

最新回复(0)