暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

【多线程 】ArrayBlockingQueue源码分析

卡布奇诺海晨 2023-08-07
97

前言

    java.util.concurrent包下的类专为并发编程而量身打造的,是线程安全的。而ArrayBlockingQueue便是其中的一员,那么它的线程安全底层是基于什么去实现的呢,带着疑问继续往下读。

    ArrayBlockingQueue是一种FIFO(first-in-first-out 先入先出排序)的有界阻塞队列,底层是数组(大小固定,一旦创建不可变更),支持从内部删除元素;新元素插入到队列的尾部,并且队列检索操作获取队列开头的元素。并发操作依赖于加锁的控制,支持阻塞式的入队出队操作。正因为有界,所以才会阻塞。

    队列满时,尝试put将元素放入队列将导致操作的线程发生阻塞队列空时,尝试take队列中的元素将类似地阻塞。

    加锁实现完全依赖于AQS,需要读者比较熟悉AQS 独占锁的获取过程和AQS Condition接口的实现。对ArrayBlockingQueue的源码解析,更像是了解一次AQS的最佳实践。

类间关系图

    可见他的父类还是不少的,这些也需要注意的一些细节。在并发编程里,我们常用的是BlockingQueue接口,也就是面向接口编程

核心字段

    public class ArrayBlockingQueue<Eextends AbstractQueue<E>
    implements BlockingQueue<E>, java.io.Serializable {
    // 保存队列元素的数组,一旦创建大小固定
    final Object[] items;

        // 下次出队的下标,可以理解为队头
    int takeIndex;

        // 下次入队的下标,可以理解为队尾
    int putIndex;

    // 队列中元素的数量
    int count;
        // 独占可重入锁
    final ReentrantLock lock;
        // 等待取出的条件对象
    private final Condition notEmpty;
    // 等待放入的条件对象
    private final Condition notFull;

    .....此处省略n行代码.......

     }

    Condition(也称为条件队列条件变量)为一个线程提供了一种方法(await)暂停执行(去“等待”),直到被另一个线程通知(signal发送信号量/signalAll等等)。它的作用就是某种情况下阻塞线程,某条件下唤醒线程。但是它是一个接口,其实现类被封装到了AQS同步队列容器中,即它的使用与锁有关。

    构造方法

    初始化ArrayBlockingQueue时,必须指定队列的容量capacity。一个入参的构造方法,fair默认为false,即非公平锁策略。多个入参的构造方法,如果fair为true,那么在插入或删除时被阻止的线程的队列访问将按FIFO顺序处理;如果false,则访问顺序未指定。另外,可以传入集合,直接构造阻塞队列。

          /**
      * 创建一个ArrayBlockingQueue,具有给定的容量capacity和默认的非公平访问策略。
           * 且容量不能小于等于0,不然抛异常
      */
      public ArrayBlockingQueue(int capacity) {
      // 调用重载的构造方法
      this(capacity, false);
      }

          /**
      * 创建一个{@code ArrayBlockingQueue},具有给定的(已修复)容量和指定的访问策略。     


      */
      public ArrayBlockingQueue(int capacity, boolean fair) {
      // 容量不能<=0,否则抛异常
      if (capacity <= 0)
      throw new IllegalArgumentException();
      // 指定容量,初始化数组
      this.items = new Object[capacity];
      // 初始化可重入锁
      lock = new ReentrantLock(fair);
      // 初始化条件变量
      notEmpty = lock.newCondition();
      notFull = lock.newCondition();
      }


          /**
      * 可传入带一定容量元素的集合遍历添加到阻塞队列,集合必须有元素否则报空指针异常
           * 
      */
      public ArrayBlockingQueue(int capacity, boolean fair,
      Collection<? extends E> c) {
      // 调用上面重载的构造方法
      this(capacity, fair);


      final ReentrantLock lock = this.lock;
      // 加锁仅用于可见性,而非互斥
      lock.lock();
      int i = 0;
      try {
      // 遍历赋值,每个元素都不能为空
      for (E e : c) {
      checkNotNull(e);
      items[i++] = e;
      }
      // 如果传入集合的个数超过了容量capacity,抛出异常被catch,再报具体异常入参不合法
      } catch (ArrayIndexOutOfBoundsException ex) {
      throw new IllegalArgumentException();
      }
      // 循环i次赋值,队列元素就有i个
      count = i;
                  // 边界处理,计算下一个元素的存放下标,如果等于容量,重置为0;否则为i
      putIndex = (i == capacity) ? 0 : i;
      } finally {
      lock.unlock();
      }
      }

      入队和出队统一逻辑

      下面两个方法,都是在持有锁的情况下被调用

      入队enqueue

            /**
        * 这个方法封装了统一的入队操作
        */
        private void enqueue(E x) {
        // assert lock.getHoldCount() == 1;
        // assert items[putIndex] == null;
        final Object[] items = this.items;
        // 在【putIndex】当前放置位置插入元素【x】到数组
        items[putIndex] = x;
        // 计算下一个元素存放的下标位置,边界处理重置为0
        if (++putIndex == items.length)
        putIndex = 0;
        // 队列中元素的数量+1
        count++;
        // 既然已入队,那么队列就是非空状态,阻塞等待获取元素的线程可以被唤醒获取元素了
        notEmpty.signal();
        }

        1、在【putIndex】当前放置位置插入元素【x】到数组

        2、计算下一个元素存放的下标位置,边界处理重置为0

        3、由于该方法是在持有锁的情况下被调用的,故元素个数递增,值都是从主内存中获取,不会存在不可见性问题,而且更新也会立即刷新回内存。

        4、既然已入队,那么队列就是非空状态,notEmpty条件变量下因take操作而阻塞等待获取元素的线程可以被唤醒获取元素了

        出队dequeue

              private E dequeue() {
          // assert lock.getHoldCount() == 1;
          // assert items[takeIndex] != null;
          final Object[] items = this.items;
          @SuppressWarnings("unchecked")
          // 取出takeIndex下标的元素
          E x = (E) items[takeIndex];
          // 设置为null
          items[takeIndex] = null;
          // 从新设置头下标
          if (++takeIndex == items.length)
          takeIndex = 0;
          // 队列元素-1
          count--;
          // 更新迭代器中的元素,itrs只有在使用迭代器的时候才会实例化
          if (itrs != null)
          itrs.elementDequeued();
          // 唤醒notFull的条件队列因调用put入队操作而被阻塞的一个线程
          notFull.signal();
          return x;
          }

          入队

          add入队,满抛异常

            // ArrayBlockingQueue.java
            public boolean add(E e) {
            // 实际上是会调用自己的offer方法入队,只不过多了队满抛异常
            return super.add(e);
            }

            // AbstractQueue.java
            public boolean add(E e) {
            if (offer(e))
            return true;
            else
                        // 队满抛异常
            throw new IllegalStateException("Queue full");
            }

            // Queue.java(接口文件)
            boolean offer(E e);

            add
            的实现是依靠父类的add
            实现,后者又依靠于子类的offer
            实现入队。所以,add
            就是在调用自己的offer
            方法,只不过有点绕,队满处理方式不同,抛异常处理。

            offer入队,返回结果true/false

                  public boolean offer(E e) {
              checkNotNull(e);
              final ReentrantLock lock = this.lock;
              // 获取锁
              lock.lock();
              try {
              // 如果当前容纳的元素个数已经等于数组长度,那么返回false
              if (count == items.length)
              return false;
              else {
              // 队列未满,将元素插入到队列中,返回true
              enqueue(e);
              return true;
              }
              } finally {
              // 释放锁
              lock.unlock();
              }
              }
              • 入队是一个写操作,自然需要加锁。lock.lock()
                不响应中断。

              • 队列已满,则无法入队,返回false

              • 队列未满,则可以入队,返回true。

              put入队,在队列的尾部插入指定的元素,如果队列已满,则阻塞等待直到空间变为可用。

                    public void put(E e) throws InterruptedException {
                // 非空校验
                checkNotNull(e);
                final ReentrantLock lock = this.lock;
                // 可中断的获取锁
                lock.lockInterruptibly();
                try {
                // 当线程从等待中被唤醒时,会比较当前队列是否已经满了
                while (count == items.length)
                notFull.await();
                // 入队
                enqueue(e);
                } finally {
                // 释放锁
                lock.unlock();
                }
                }

                在进入加锁代码之前,执行的是lock.lockInterruptibly()。这意味着,当前线程在抢到锁之前,如果被中断了,put方法会抛出中断异常。

                进入加锁代码之后,当前线程便已是获得了锁。

                如果队列未满,那么根本不会执行notFull.await(),直接入队。

                需要使用while (count == items.length)来防止虚假唤醒,即使当前线程从notFull.await()恢复执行了,如果当前队列还是满的,那么应该重新进入条件队列。所以,需要重新检查一遍count == items.length。

                你可能会产生疑问,为什么需要重新检查一遍。因为当前线程从notFull.await()恢复执行,一定是因为别的线程执行了notFull.signal()(别的线程的这个时间点,队列确实未满)。但由于当前线程是从AQS的条件队列(等待)转移到AQS的同步队列队尾(参与锁的争夺),而排在同步队列前面的其他线程也有可能去执行入队操作,可能等到当前线程获得锁后(所以才会从notFull.await()恢复执行),队列又变成满了。

                此put函数只有成功入队后,才可能从put调用处返回。

                当队列未满,则入队。

                offer入队+超时机制

                      public boolean offer(E e, long timeout, TimeUnit unit)
                  throws InterruptedException {

                  checkNotNull(e);
                  long nanos = unit.toNanos(timeout);
                  final ReentrantLock lock = this.lock;
                  lock.lockInterruptibly();
                  try {
                  while (count == items.length) {
                  // 如果队列是满的,且等待时间<= 0这代表超时,所以直接返回false
                  if (nanos <= 0)
                  return false;
                  nanos = notFull.awaitNanos(nanos);
                  }
                  enqueue(e);
                  return true;
                  } finally {
                  lock.unlock();
                  }
                  }

                  相比上一个实现,使用的是awaitNanos。

                  从notFull.awaitNanos(nanos)返回有三种原因:超时前的signal、超时前的中断、超时。

                  超时前的signal。只有这种情况,才可能返回一个大于0的数字。

                  超时前的中断。返回时,抛出中断异常。

                  超时(不管之后有没有中断)。只可能返回一个小于0的数字。

                  因为超时前的signal而从notFull.awaitNanos(nanos)返回,需要进行虚假唤醒的检查。如果此时队列还是满的,当前线程再次进入AQS的条件队列;如果此时队列确实未满,那么入队,返回true。

                  如果此时队列是满的,当前线程再次进入AQS的条件队列之前,需要检查剩余时间是否大于0,如果不是大于0,说明在awaitNanos上花费的时间已经超过了限制,则返回false。

                  某种情景再现:

                  当前线程调用notFull.awaitNanos(500),准备进行500ns的等待。

                  别的线程在剩余时间大约还有300ns的时间时,调用了notFull.signal(),唤醒了当前线程。

                  当前线程从notFull.awaitNanos(500)处返回,返回值为300。

                  循环继续,检查却发现队列已满。

                  if (nanos <= 0)不满足,继续执行notFull.awaitNanos(300)。

                  当前线程继续等待300ns。

                  总结

                  出队

                  peek出队

                        public E peek() {
                    final ReentrantLock lock = this.lock;
                    lock.lock();
                    try {
                    return itemAt(takeIndex); // null when queue is empty
                    } finally {
                    lock.unlock();
                    }
                    }

                    直接返回索引处元素,可能为null(队列为空),正如peek的含义,只获取不出队。

                    poll出队

                          public E poll() {
                      final ReentrantLock lock = this.lock;
                      lock.lock();
                      try {
                      return (count == 0) ? null : dequeue();
                      } finally {
                      lock.unlock();
                      }
                      }

                      此函数可能返回null,当队列为空时。

                      take出队,队列空会阻塞等待

                            public E take() throws InterruptedException {
                        final ReentrantLock lock = this.lock;
                        lock.lockInterruptibly();
                        try {
                        // 虚假唤醒检查
                        while (count == 0)
                        notEmpty.await();
                                    // 如果队列确实不空,那么执行出队动作
                        return dequeue();
                        } finally {
                        lock.unlock();
                        }
                        }

                        poll出队+超时机制

                              public E poll(long timeout, TimeUnit unit) throws InterruptedException {
                          long nanos = unit.toNanos(timeout);
                          final ReentrantLock lock = this.lock;
                          lock.lockInterruptibly();
                          try {
                          // 虚假唤醒检查
                          while (count == 0) {
                          // 如果队列是空的,且等待时间<= 0这代表不用等待,所以直接返回null
                          if (nanos <= 0)
                          return null;
                          nanos = notEmpty.awaitNanos(nanos);
                          }
                          // 如果队列确实不空,那么执行出队动作
                          return dequeue();
                          } finally {
                          lock.unlock();
                          }
                          }

                          总结

                          删除

                          remove该函数如果删除的不是队首元素,会涉及到整体移动的过程,可能会比较耗时,不建议使用。

                          现在队列中非null元素的范围是[takeIndex, putIndex)
                          的左闭右开的区间。

                                public boolean remove(Object o) {
                            if (o == null) return false;
                            final Object[] items = this.items;
                            final ReentrantLock lock = this.lock;
                            lock.lock();
                            try {
                            if (count > 0) {//队列有元素存在
                            final int putIndex = this.putIndex;
                            int i = takeIndex;
                            do {
                            if (o.equals(items[i])) {
                            removeAt(i);
                            return true;
                            }
                            if (++i == items.length)
                            i = 0;
                            } while (i != putIndex);
                            //到达区间[takeIndex, putIndex)的边界,说明所有非null元素都找遍了
                            }
                            return false;//没有找到元素
                            } finally {
                            lock.unlock();
                            }
                            }

                            循环从[takeIndex, putIndex)
                            的左边界开始,直到右边界结束。如果找到元素,则删除它。

                                  void removeAt(final int removeIndex) {
                              // assert lock.getHoldCount() == 1;
                              // assert items[removeIndex] != null;
                              // assert removeIndex >= 0 && removeIndex < items.length;
                              final Object[] items = this.items;
                              if (removeIndex == takeIndex) {//如果刚好删除的是队首,那刚好是一个出队动作
                              // removing front item; just advance
                              items[takeIndex] = null;
                              if (++takeIndex == items.length)
                              takeIndex = 0;
                              count--;
                              if (itrs != null)
                              itrs.elementDequeued();
                              } else {//其他情况
                              //[i,putIndex)区间内的第一个元素被删除,需要往左压实这个区间
                              final int putIndex = this.putIndex;
                              for (int i = removeIndex;;) {
                              int next = i + 1;
                              if (next == items.length)
                              next = 0;
                              if (next != putIndex) {//还没到达边界
                              items[i] = items[next];//将后面的复制到前面去
                              i = next;
                              } else {//到达边界
                              items[i] = null;//清空区间内最后一个元素
                              this.putIndex = i;//最后putIndex当然也得左移,i此时肯定是putIndex - 1
                              break;
                              }
                              }
                              count--;
                              if (itrs != null)
                              itrs.removedAt(removeIndex);
                              }
                              notFull.signal();
                              }

                              • 如果刚好删除的队首元素,那刚好是一次出队操作。

                              • 如果是其他情况,现在删除的是i
                                索引元素,但为了队列非null元素连续(考虑循环数组也得连续),那么[i, putIndex)
                                区间内的第一个元素已经被删除变成null了,需要往左压实,即[i+1, putIndex)
                                内的元素整体左移。

                              总结

                              • 当队列为空或为满时,takeIndex putIndex二者才会相同。

                              • 所有常用操作都需要加锁,甚至是属于读操作的peek,因为加锁强制内存刷新,能让线程看到最新的队列。

                              • 入队出队操作,都有一次尝试版本,和阻塞等待版本。

                              • 使用Lock来控制并发操作。

                              • 两个Condition的使用,是控制阻塞等待的关键。

                              • 删除操作支持删除内部元素。

                              • 入队出队都是同一把锁,锁的机制是ReentrantLock+Condition


                              文章转载自卡布奇诺海晨,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

                              评论