资讯动态

Java 阻塞队列、延迟队列

发布时间:2026/8/5 14:53:41 来源:尧图企业网站定制
目录一、阻塞队列1.LinkedBlockingQueue1.1 结构1.2 put()方法1.3 take()方法2.ArrayBlockingQueue2.1 结构2.2 put()方法2.3 take()方法3.面试题实现生产者消费者模型二、延迟队列1.DelayQueue1.1 结构1.2 offer()方法1.3 take()方法2.面试题实现订单超时取消了解该类的原理需要按照顺序学习前置知识CAS-Unsafe-AQS-ReentrantLock一、阻塞队列阻塞队列就是用于存取数据的共享队列基于ReentrantLock保证多线程操作队列的线程安全。ArrayBlockingQueue共享队列采用环形数组实现使用单ReentrantLock保证所有生产者消费者串行访问数据队列。LinkedBlockingQueue共享队列采用双向链表实现使用双ReentrantLock保证生产者之间串行、消费者之间串行访问数据队列。1.LinkedBlockingQueue1.1 结构publicclassLinkedBlockingQueueEextendsAbstractQueueE{// 队列容量privatefinalintcapacity;// 当前队列元素数量这里用AtomicInteger是因为生产者和消费者不是串行操作count所以不能只用volatileprivatefinalAtomicIntegercountnewAtomicInteger();// 头节点transientNodeEhead;// 尾节点privatetransientNodeElast;// 保证消费者线程串行privatefinalReentrantLocktakeLocknewReentrantLock();// 链表空了消费者等待队列privatefinalConditionnotEmptytakeLock.newCondition();// 保证生产者线程串行privatefinalReentrantLockputLocknewReentrantLock();// 链表满了生产者等待队列privatefinalConditionnotFullputLock.newCondition();}其中head、last作为首尾结点形成双向链表链表用于记录消息takeLock控制消费者之间串行从链表头取消息putLock控制生产者之间串行向链表尾存消息。当链表为空时消费者取消息就要进入notEmpty等待当链表存满时生产者存消息就要进入notFull等待。capacity表示容量防止链表容量无线增长count记录链表当前长度。1.2 put()方法生产者调用put()向链表中存消息步骤如下putLock.lock()获取锁保证多生产者串行向表尾存消息。-如果count.get() capacity表示链表存满了notFull.await()进入等待队列等待。后续被唤醒跳到第1步。-如果链表还有容量插入新结点到表尾count.getAndIncrement()更新容量。如果链表还有容量notFull.signal()唤醒一个生产者。putLock.unlock()释放锁。如果插入结点前链表是空的那么插入结点后链表有消息了takeLock.lock()获取锁notEmpty.signal()唤醒一个消费者takeLock.unlock()释放锁。publicvoidput(Ee)throwsInterruptedException{if(enull)thrownewNullPointerException();// 记录链表put前的容量intc-1;NodeEnodenewNodeE(e);finalReentrantLockputLockthis.putLock;finalAtomicIntegercountthis.count;// 生产者获取锁putLock.lockInterruptibly();try{// 队列满了进入等待状态while(count.get()capacity){notFull.await();}// 队列不满存消息enqueue(node);// 更新当前容量返回更新前的容量ccount.getAndIncrement();// 如果还有容量那么唤醒另一个生产者if(c1capacity)notFull.signal();}finally{// 释放锁putLock.unlock();}// 如果队列之前是空的if(c0){finalReentrantLocktakeLockthis.takeLock;// 获取锁takeLock.lock();try{// 唤醒消费者notEmpty.signal();}finally{// 释放锁takeLock.unlock();}}}1.3 take()方法消费者调用take()从链表中取消息步骤如下takeLock.lock()获取锁保证多消费者串行从表头取消息。-如果count.get() 0表示链表为空notEmpty.await()进入等待队列等待。后续被唤醒跳到第1步。-如果链表不为空从表头取消息count.getAndIncrement()更新容量。如果链表还不为空notEmpty.signal()唤醒一个生产者。takeLock.unlock()释放锁。如果取出结点前链表是满的那么取出结点后链表就有容量了putLock.lock()获取锁notFull.signal()唤醒一个生产者putLock.unlock()释放锁。publicEtake()throwsInterruptedException{Ex;// 记录链表take前的容量intc-1;finalAtomicIntegercountthis.count;finalReentrantLocktakeLockthis.takeLock;// 消费者获取锁takeLock.lockInterruptibly();try{while(count.get()0){// 队列为空进入等待状态notEmpty.await();}// 队列不空取消息xdequeue();// 更新当前容量返回更新前的容量ccount.getAndDecrement();// 如果还有消息那么唤醒另一个消费者if(c1)notEmpty.signal();}finally{// 释放锁takeLock.unlock();}// 如果队列之前是满的if(ccapacity){finalReentrantLocktakeLockthis.takeLock;// 获取锁takeLock.lock();try{// 唤醒生产者notEmpty.signal();}finally{// 释放锁takeLock.unlock();}}}returnx;}2.ArrayBlockingQueueArrayBlockingQueue与LinkedBlockingQueue的唯一区别是所有生产者和消费者共用同一把ReentrantLock不仅生产者之间串行消费者之间串行、生产者和消费者之间也是串行的。2.1 结构publicclassArrayBlockingQueueEextendsAbstractQueueEimplementsBlockingQueueE{// 存储消息的数组finalObject[]items;// 取消息索引inttakeIndex;// 存消息索引intputIndex;// 当前队列元素数量因为所有线程共用一把锁count操作都在临界区所以不用加volatileintcount;// 所有线程串行操作环形数组finalReentrantLocklock;// 数组空了消费者等待队列privatefinalConditionnotEmpty;// 数组满了生产者等待队列privatefinalConditionnotFull;}其中lock控制生产者和消费者串行从环形数组items存取消息putIndex、takeIndex记录生产者、消费者存取元素索引。当数组为空时消费者取消息就要进入notEmpty等待当数组存满时生产者存消息就要进入notFull等待。count记录数组当前元素数量。2.2 put()方法生产者调用put()向数组中存消息步骤如下lock.lock()获取锁保证多生产者串行向数组存消息。-如果count item.length表示数组存满了notFull.await()进入等待队列等待。后续被唤醒跳到第1步。-如果数组还有容量插入新结点count更新容量。notEmpty.signal()唤醒一个消费者。lock.unlock()释放锁。publicvoidput(Ee)throwsInterruptedException{//确保插入的元素不为nullcheckNotNull(e);finalReentrantLocklockthis.lock;// 生产者获取锁lock.lock();try{// 队列满了进入等待状态while(countitems.length)notFull.await();finalObject[]itemsthis.items;// 队列不满存消息items[putIndex]e;// 更新下一个存消息的索引if(putIndexitems.length)putIndex0;// 数组消息数1count;// 唤醒消费者notEmpty.signal();}finally{// 释放锁lock.unlock();}}2.3 take()方法消费者调用take()从数组中取消息步骤如下lock.lock()获取锁保证多消费者串行从数组取消息。-如果count0表示数组为空notEmpty.await()进入等待队列等待。后续被唤醒跳到第1步。-如果数组不为空从数组取消息count--更新容量。notFull.signal()唤醒一个生产者。lock.unlock()释放锁。publicEtake()throwsInterruptedException{finalReentrantLocklockthis.lock;// 消费者获取锁lock.lock();try{// 队列为空进入等待状态while(count0)notEmpty.await();finalObject[]itemsthis.items;// 队列不空取消息Ex(E)items[takeIndex];// 将takeIndex中的消息移出数组items[takeIndex]null;// 更新下一个存消息的索引if(takeIndexitems.length)takeIndex0;// 数组消息数-1count--;if(itrs!null)itrs.elementDequeued();// 唤醒生产者notFull.signal();returnx;}finally{// 释放锁lock.unlock();}}3.面试题实现生产者消费者模型阻塞队列不仅要管理读线程、还要管理写线程、所以需要两个ReentrantLock一个控制生产者串行存入任务一个控制消费者串行取出注意不是执行任务然后还需要两个LockSupport对象控制生产者消费者之间的同步这些所有过程封装为以下两个方法put(E e)将元素插入队列中如果队列已满该方法会一直阻塞直到队列有空间可用或者线程被中断。take()获取并移除队列头部的元素如果队列为空该方法会一直阻塞直到队列非空或者线程被中断。classCookerextendsThread{ArrayBlockingQueueStringblockingQueue;publicCooker(ArrayBlockingQueueStringblockingQueue){this.blockingQueueblockingQueue;}Overridepublicvoidrun(){while(true){blockingQueue.put(food);}}}classCustomerextendsThread{ArrayBlockingQueueStringblockingQueue;publicCustomer(ArrayBlockingQueueStringblockingQueue){this.blockingQueueblockingQueue;}Overridepublicvoidrun(){while(true){blockingQueue.take();}}}publicclassTestThread{publicstaticvoidmain(Stringargs[]){// main线程ArrayBlockingQueueStringblockingQueuenewArrayBlockingQueue(10);CookercookernewCooker(blockingQueue);CustomercustomernewCustomer(blockingQueue);cooker.start();customer.start();}}二、延迟队列DelayQueue基于PriorityQueue小顶堆ReentrantLock实现用于实现延迟任务比如超时订单取消并退还库存。只有getDelay()0即元素过期时才能从队列中取出元素。DelayQueue中存放的任务必须实现Delayed接口并且需要重写getDelay()方法给定堆排序的规则publicinterfaceDelayedextendsComparableDelayed{longgetDelay(TimeUnitunit);}1.DelayQueue1.1 结构publicclassDelayQueueEextendsDelayedextendsAbstractQueueEimplementsBlockingQueueE{// 存放消息的堆privatefinalPriorityQueueEqnewPriorityQueueE();// 所有线程串行操作堆privatefinaltransientReentrantLocklocknewReentrantLock();// 第一个等待的消费者该消费者会进入TIME_WAIT状态privateThreadleader;// 堆空了消费者等待队列privatefinalConditionavailablelock.newCondition();}其中lock控制生产者和消费者串行从堆q存取消息。当堆中没有到期的任务时消费者取消息就要进入available等待。leader记录负责消费堆头消息的消费者该消费者会处于TIME_WAIT状态。延迟队列是无界的PriorityQueue存满后会自动扩容。1.2 offer()方法生产者调用offer()向堆中存消息步骤如下lock.lock()获取锁保证多生产者串行向堆存消息。q.offer(e)存入新消息更新堆排列。如果当前消息位于堆头q.peek()e说明该消息是堆中最先到期的消息将leader置空并available.signal()唤醒队头消费者。lock.unlock()释放锁。publicbooleanoffer(Ee){finalReentrantLocklockthis.lock;// 生产者获取锁lock.lock();try{// 将消息存放到堆中q.offer(e);// 本次入队的消息位于队头if(q.peek()e){// 将leader设置为空leadernull;// 唤醒消费者available.signal();}returntrue;}finally{// 释放锁lock.unlock();}}1.3 take()方法消费者调用take()从堆中取消息步骤如下lock.lock()获取锁保证多消费者串行从堆中取消息。-如果q.peek()null表示堆为空available.await()进入等待队列等待。后续被唤醒跳到第1步。如果堆不为空获取堆头消息检查消息是否到期Delayed.getDelay()如果到期取出消息并更新堆排序。-如果堆不为空堆头消息没到期且leader!null说明堆头元素已经有消费者在等待了available.await()进入等待队列等待。后续被唤醒跳到第1步。-如果堆不为空堆头消息没到期且leadernull说明还没有为堆头元素分配消费者那么更新leaderThread.currentThread()为当前线程并available.awaitNanos(delay)进入TIME_WAIT状态。后续被唤醒跳到第1步。取出消息后新的堆头消息可能还没有分配消费者available.signal()唤醒一个消费者。takeLock.lock()释放锁。publicEtake()throwsInterruptedException{finalReentrantLocklockthis.lock;// 消费者获取锁lock.lockInterruptibly();try{for(;;){// 查看堆头消息Efirstq.peek();// 堆为空进入等待状态if(firstnull)available.await();else{// 获取堆头消息到期时间longdelayfirst.getDelay(NANOSECONDS);// 堆头消息已到期取出元素if(delay0)returnq.poll();firstnull;// 堆不为空且leader不为空说明已经分配了leader线程进入TIME_WAIT状态等待消费堆头消息当前线程进入等待队列if(leader!null)available.await();else{// 堆不为空且leader为空说明还没有为堆头消息分配消费线程ThreadthisThreadThread.currentThread();leaderthisThread;try{// 当前消费者进入有TIME_WAIT有限等待available.awaitNanos(delay);}finally{// 当前线程取出了堆头消息新的堆头消息还没有分配消费者将leader置空if(leaderthisThread)leadernull;}}}}}finally{// 当leader为null并且堆中有任务时唤醒一个消费者if(leadernullq.peek()!null)available.signal();// 释放锁lock.unlock();}}答疑7. 延迟队列默认是公平模式吗默认非公平模式消费者获取锁后检查堆头元素是否到期到期就会取出消息不会直接进入等待队列。8.leader引用的作用其实延迟队列完全可以照搬阻塞队列唯一区别就是将WAIT状态变成TIME_WAIT状态但是这么设计有个问题堆头消息可能被所有线程绑定因此每当堆头消息到期就会出现惊群效应本质是消费者会根据堆头消息到期时间进入TIME_WAIT状态。DelayQueue的思路是用leader记录堆头消息是否有消费者绑定了只要有绑定那么其他消费者就进入无限等待WAIT状态解决了惊群问题。9. 堆头消息在等待队列中只有一个消费者等待吗不是如果堆头消息已经绑定了消费者此时新的消息到期时间更短成为了新的堆头消息如果该消息到期被取出那么leader就会被置空标识堆头消息还没有指定消费者但实际上堆头消息已经有消费者等待了就会出现等待队列中有多个消费者等待堆头消息但是相比于惊群效应少了很多本质是leader作为延迟队列的成员变量而不是消息的成员变量如果作为消息的成员变量那么就是一对一等待关系。2.面试题实现订单超时取消publicclassMain{staticclassOrderTaskimplementsDelayed{privateStringorderId;// 订单IDprivatelongexpireTime;// 到期时间publicOrderTask(StringorderId,longdelayTime){this.orderIdorderId;this.expireTimeSystem.nanoTime()TimeUnit.MILLISECONDS.toNanos(delayTime);}// 获取剩余延迟时间OverridepubliclonggetDelay(TimeUnitunit){returnunit.convert(expireTime-System.nanoTime(),TimeUnit.NANOSECONDS);}// 小顶堆排序OverridepublicintcompareTo(Delayedother){OrderTasko(OrderTask)other;returnLong.compare(this.expireTime,o.expireTime);}}publicstaticvoidmain(String[]args){DelayQueueOrderTaskdelayQueuenewDelayQueue();newThread(()-{delayQueue.offer(newOrderTask(1,5000));delayQueue.offer(newOrderTask(2,1000));delayQueue.offer(newOrderTask(3,6000));}).start();try{while(true)System.out.println(订单已被处理delayQueue.take().orderId);}catch(InterruptedExceptione){e.printStackTrace();}}}

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价