ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Java队列从基础到实战:从Queue到BlockingQueue再到消息队列避坑指南

Java队列从基础到实战:从Queue到BlockingQueue再到消息队列避坑指南 1. 从一次线上事故聊起队列不只是“先进先出”1.1 那次线程池满了我才开始认真研究队列大概是在三年前我负责的一个批处理系统在618大促期间突然OOM了。现象很典型定时任务每五分钟扫描一次待处理订单扫描量一旦上来机器内存就一路飙升然后监控告警轰炸服务重启后才能恢复但过不了半小时又挂一次。当时第一反应是下游服务响应慢、数据库连接池不够于是加超时、加熔断、加机器全都没用。最后抓了一把堆转储才找到元凶线程池用的是Executors.newFixedThreadPool(10)而它底层默认的队列是一个无界的LinkedBlockingQueue。任务一旦积压队列可以无限增长内存压垮只是时间问题。这事给我上了一课队列不是简单的“先进先出”在Java里它横跨了集合框架、并发包、线程池、消息中间件好几个层面。选错队列轻则性能差重则线上事故。所以今天我想把Java里的队列好好梳理一遍从最基础的Queue接口一直到生产环境里的消息队列重复消费问题结合我踩过的坑把该讲清楚的都讲清楚。1.2 队列在Java里的三个层次很多刚入行的同学以为队列就是LinkedList或者ArrayBlockingQueue其实在Java生态里“队列”至少有三个层次层次代表解决什么问题集合框架层Queue、Deque、PriorityQueue普通的内存数据结构单线程或多线程无锁并发时使用并发包层BlockingQueue系列线程间数据传递、生产者消费者模型、线程池任务队列中间件层RocketMQ、Kafka、RabbitMQ跨进程、跨服务、分布式环境下的消息解耦与削峰这三个层次并不是递进替代关系而是对应不同场景。比如你写个算法题、处理一段内存数据用ArrayDeque就行你要在线程之间传递任务就必须用BlockingQueue你要把订单变更通知发给下游多个服务那得靠MQ。下面我从易到难把每一层的核心知识点和实战经验拆开说。2. 集合框架里的Queue和Deque选型与陷阱2.1 别再迷信LinkedList了默认用ArrayDeque搜索热词里你肯定经常看到“队列”和“栈”并列出现Java里LinkedList确实同时实现了List和Deque接口所以很多教程会用LinkedList来演示队列操作QueueString queue new LinkedList(); queue.offer(订单A); queue.offer(订单B); System.out.println(queue.poll()); // 订单A这个写法没错但在性能和内存上并不理想。LinkedList的每个节点都是一个独立的Node对象需要额外的指针和对象头而且匹配不连续CPU缓存命中率低。如果只做队列和栈的标准操作完全可以用ArrayDeque代替ArrayDequeString queue new ArrayDeque(); queue.offer(订单A); queue.offer(订单B); System.out.println(queue.poll());ArrayDeque底层是一个循环数组加减头尾元素都是O(1)还省掉了大量节点对象的开销。所以我的习惯是只要不需要中间插入或者不需要按下标访问一律用ArrayDeque而不是LinkedList。但是有个前提要注意ArrayDeque不允许null元素而且它不是线程安全的。如果你在多个线程里同时offer和poll需要自己加锁或者直接上并发的BlockingQueue。2.2 PriorityQueue排序规则搞错就是灾难PriorityQueue也是队列家族的一员但它不保证先进先出而是按优先级弹出元素底层是二叉堆。PriorityQueueInteger priorityQueue new PriorityQueue(); priorityQueue.offer(5); priorityQueue.offer(1); priorityQueue.offer(3); System.out.println(priorityQueue.poll()); // 1 priorityQueue.offer(2); System.out.println(priorityQueue.poll()); // 2 System.out.println(priorityQueue.poll()); // 3默认情况下它是最小堆也就是每次弹出最小值。如果想弹出最大值要传入自定义比较器PriorityQueueInteger maxHeap new PriorityQueue((a, b) - b - a);这里面有两个经典坑迭代顺序不等于优先级顺序。PriorityQueue的iterator()遍历结果是无序的只有poll()和peek()才能按堆的性质拿到极值。所以千万不要用toArray()然后宣称已经排序。自定义对象的比较器必须与equals保持一致。比如你用PriorityQueue存储订单按金额排序同时equals按订单ID判断同一订单这时候remove方法可能找不到你要删的元素。因为remove先通过equals找元素而堆结构调整依赖compareTo/Comparator两者不一致会导致堆性质被破坏。2.3 双端队列如何实现滑动窗口热搜词里有“数据结构 双端队列”“单调队列-滑动窗口”这些在面试和算法竞赛里都很常考。Deque的价值就在这里体现它支持在头尾两端高效插入和删除正好满足单调队列的维护需要。以“求滑动窗口最大值”为例核心思路是用双端队列存储元素的下标并保证队列内元素对应值单调递减。每次窗口移动队首元素如果不在窗口内就pollFirst()出队新元素入队前把队尾所有比它小的元素都pollLast()掉因为它们不可能再成为窗口内的最大值了当前元素入队尾队首就是窗口最大值。public int[] maxSlidingWindow(int[] nums, int k) { int n nums.length; int[] result new int[n - k 1]; ArrayDequeInteger deque new ArrayDeque(); for (int i 0; i n; i) { // 移除不在窗口内的队首下标 while (!deque.isEmpty() deque.peekFirst() i - k 1) { deque.pollFirst(); } // 移除队尾所有比当前值小的元素 while (!deque.isEmpty() nums[deque.peekLast()] nums[i]) { deque.pollLast(); } deque.offerLast(i); if (i k - 1) { result[i - k 1] nums[deque.peekFirst()]; } } return result; }这个算法的均摊复杂度是O(n)比用优先队列 O(n log n) 更优雅。所以下次看到滑动窗口最值问题优先考虑双端队列实现单调队列。3. 并发场景下BlockingQueue才是真正的核心3.1 阻塞与非阻塞的本质在java.util.concurrent包出现之前线程间传数据要么用wait/notify要么用轮询加锁。轮询浪费CPUwait/notify又容易写错。BlockingQueue把这块逻辑封装好了它的put方法在队列满时会自动阻塞当前线程直到队列有空位take方法在队列空时会自动阻塞直到有元素进来。这里面的底层机制其实是锁和条件变量。以ArrayBlockingQueue为例它内部维护了一个ReentrantLock和两个ConditionnotEmpty、notFull。当队列满时put就调用notFull.await()当队列空时take就调用notEmpty.await()。线程被park住不会空转消耗CPU。非阻塞的对应物是ConcurrentLinkedQueue它基于CAS实现没有锁但也没有阻塞能力。你不能指望一个线程在队列空时“等一下”因为poll()会直接返回null。所以在生产者消费者模型、线程池任务队列场景中BlockingQueue是首选。3.2 几款常用实现怎么选我把常用的阻塞队列列成一张对比表方便直接参照实现类是否有界数据结构锁策略典型场景ArrayBlockingQueue有界容量必须指定循环数组一把锁 两个Condition可配公平锁有界缓冲、线程池任务队列LinkedBlockingQueue默认无界容量为Integer.MAX_VALUE链表两个锁takeLock和putLock线程池默认任务队列、无界缓冲SynchronousQueue无缓冲无内部容量基于CAS/锁直接交接线程池对应CachedThreadPoolLinkedTransferQueue无界链表无锁CAS高并发短任务传递支持transferPriorityBlockingQueue无界二叉堆一把锁优先级任务调度DelayQueue无界基于PriorityQueue一把锁延迟任务、定时器很多人选型时只看“是不是线程安全”忽略了“是否有界”这个致命属性。LinkedBlockingQueue如果不指定容量它就是一个可以无限增长的队列任务积压时内存会持续上涨前面提到的OOM事故就是这么来的。3.3 手写一个生产者消费者模型网上到处都是生产者消费者例子但写得规范的不多。我分享一个实际会用到的写法核心是两条用poll/offer的带超时版本而不是死等put/take在线程池的submit之外要能优雅关闭。public class OrderTaskProcessor { private final BlockingQueueString queue; private final ExecutorService consumerExecutor; private volatile boolean running true; public OrderTaskProcessor(int queueCapacity, int consumerCount) { this.queue new ArrayBlockingQueue(queueCapacity); this.consumerExecutor Executors.newFixedThreadPool(consumerCount); for (int i 0; i consumerCount; i) { consumerExecutor.submit(this::consumeLoop); } } public boolean submitOrder(String orderId) { // 最多阻塞200ms避免生产者被打满 try { return queue.offer(orderId, 200, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return false; } } private void consumeLoop() { while (running) { try { String orderId queue.poll(5, TimeUnit.SECONDS); if (orderId ! null) { handleOrder(orderId); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); if (running) { // 恢复消费 continue; } break; } } } private void handleOrder(String orderId) { // 实际处理逻辑 System.out.println(处理订单: orderId); } public void shutdown() { running false; consumerExecutor.shutdown(); try { consumerExecutor.awaitTermination(30, TimeUnit.SECONDS); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }这个例子有几个关键点用了ArrayBlockingQueue并指定容量防止无界堆积offer带了超时时间生产者不会无限阻塞poll也带了超时时间消费者线程在队列空时不会一直空转也不会在runningfalse时因为take()卡住退出不了通过volatile running实现优雅关闭比直接shutdownNow更可控。3.4 使用BlockingQueue的几个坑第一个坑自定义比较器的优先级队列不能保证全局有序之前提过。第二个坑add、remove、element这几个方法在队列满/空时直接抛异常而offer、poll、peek返回特殊值put、take会阻塞。搞清楚三组方法的语义可以避免很多线上异常操作失败时行为使用建议add(e)队列满抛IllegalStateException不要用于业务流控offer(e)返回false用于环保入队put(e)阻塞直到有空间用于生产者必须送达的场景remove()空队列抛NoSuchElementException不要用于消费者poll()返回null用于非阻塞消费take()阻塞直到有元素用于消费者必须等待第三个坑中断异常必须正确处理。take和put在等待时如果线程被中断会抛出InterruptedException。最差的做法是打印日志后继续吞掉异常正确的做法是重新设置中断标志位Thread.currentThread().interrupt()方便上层感知。4. 线程池的队列选择最容易埋雷的一个配置4.1 线程池和队列的协作流程面试题里“线程池参数”几乎是必考但很多人只背了corePoolSize、maxPoolSize忽略了workQueue才是线程池暴雷的根源。ThreadPoolExecutor任务提交流程是这样的当前工作线程数小于corePoolSize创建核心线程执行核心线程满了新任务放进workQueue排队队列满了再创建非核心线程直到maxPoolSize再满执行拒绝策略。关键就在第二步队列不是用来无限装任务的而是用来缓冲峰值压力的。如果队列无边第3步和第4步根本不会生效线程池最多只有核心线程数所有任务都堆在内存队列里。以FixedThreadPool为例它用的就是无界LinkedBlockingQueue。你能保证所有任务都能在下游恢复之前被消化完吗不能。所以高并发生产环境里我几乎不用Executors预设的线程池而是手动创建ThreadPoolExecutor。4.2 无界队列、有界队列、同步移交怎么选直接上一张对比表结合我自己的实际经验队列类型代表实现行为特征适用场景无界队列LinkedBlockingQueue()不传容量队列容量无限拒绝策略形同虚设任务量可预测、系统压力可控的内部异步任务有界队列ArrayBlockingQueue(capacity)队列满了才创建新线程再满就拒绝需要控制内存和背压的生产消费场景同步移交SynchronousQueue几乎没有缓冲直接把任务交给线程大量短小且独立的CPU密集型任务优先级队列PriorityBlockingQueue按优先级处理任务但有界限制较弱需要任务优先级调度的场景我记得有个同事为了让接口响应快给线程池maxPoolSize设了200但workQueue用的是无界LinkedBlockingQueue。结果核心线程只有10个所有任务都往队列塞最后队列里积压了百万级任务。下游一抖动内存快速上涨线程池反而成了故障放大器。后来我们改成了ArrayBlockingQueue(1000)并把maxPoolSize也设成100配合合适的拒绝策略整个系统立刻稳了很多。4.3 拒绝策略和容量权衡拒绝策略有四个内置选项我列出真实使用率策略效果是否推荐AbortPolicy默认策略直接抛RejectedExecutionException适合能感知失败、需要补偿的流程CallerRunsPolicy让提交任务的线程自己执行任务适合慢速降级但要注意提交线程被占用DiscardPolicy静默丢弃新任务不推荐容易丢数据DiscardOldestPolicy丢弃最老的任务然后重试提交适合允许丢弃积压任务的场景我的经验是能用CallerRunsPolicy就用它宁可让调用线程卡一下也别把任务丢掉或者把内存堆爆。但前提是调用线程能承受这个额外开销。如果调用方是高频HTTP请求一旦线程池满了CallerRunsPolicy会让全部请求线程都去执行队列任务可能导致整个服务失去响应。所以更稳妥的组合是有界队列容量约等于核心线程数 × 每个任务最大容忍等待时间 ÷ 每个任务平均执行时间。比如核心线程10每个任务平均执行50ms业务可接受排队1秒那队列容量就控制在200左右。这只是估算公式真正上线前要做压测并在运行时监控队列长度。我一般会在ThreadPoolExecutor外面包一层beforeExecute/afterExecute把任务提交数、完成数、队列深度定期打进监控系统方便第一时间发现队列堆积。5. 从内置队列到消息队列重复消费是怎么来的5.1 为什么单机队列解决不了跨服务问题很多项目在刚开始时用BlockingQueue做异步化足够了。但等微服务拆分后问题就来了JVM内的队列是单机的服务重启后数据全丢而且多个实例之间无法共享队列。比如订单服务要把“支付成功”的事件发给积分服务和推送服务如果就用本地内存队列订单服务启两个实例事件只会发到其中一个实例的队列推送服务根本不知道。这才有了消息队列中间件。Kafka、RocketMQ、RabbitMQ 这些消息队列能跨进程、跨服务传递消息还提供持久化和削峰能力。但分布式环境下也带来了新的问题其中最经典的就是重复消费。5.2 消息队列重复消费的根源先看at least once这个语义消息队列为了保证不丢消息会在消费者成功消费后返回确认信息比如Kafka的commit。如果消费者在处理消息时挂掉了或者网络抖动导致确认信息没有送达这条消息就会被重新投递给消费者。举一个实际场景消费者从MQ里取出一条“给用户加100积分”的消息业务代码先执行了积分更新然后在提交offset之前服务宕机了。重启后这条消息又被拉下来积分的数据库操作又执行了一遍。用户积分的记录看起来多了100这就是重复消费。这种问题不是MQ的Bug而是分布式系统中“至少投递一次”“消费端原子性无法跨网络保证”的必然结果。想要真正解决不能指望MQ帮你过滤必须在消费端做幂等。5.3 幂等设计才是解药我把实践中验证过比较靠谱的幂等方案列一下方案思路适用场景唯一主键数据库表建唯一索引重复插入直接失败/忽略订单创建、流水记录去重表每次消费前查RedisSETNX或业务表记录消息ID通用做法状态机字段操作前检查当前状态只有期望状态才能流转订单状态更新、支付回调分布式锁对业务主键加锁释放后再判断是否处理过并发冲突高的写操作最朴素也最实用的做法是给每个消息生成一个全局唯一的messageId然后在消费方维护一张去重表。插入时把这个messageId作为唯一键如果插入冲突说明已经处理过直接跳过。public void handlePaymentMessage(PaymentMessage message) { String key order: message.getOrderId() :pay; Boolean success redisTemplate.opsForValue().setIfAbsent(key, message.getMessageId(), 24, TimeUnit.HOURS); if (success null || !success) { log.info(重复消息跳过, messageId{}, orderId{}, message.getMessageId(), message.getOrderId()); return; } // 真正的业务处理 updateOrderPaid(message.getOrderId()); }注意这里有个细节setIfAbsent加过期时间是为了避免Redis里的幂等标识无限增长。过期时间至少要覆盖可能的重复投递间隔我通常设成24小时基本覆盖了消费端宕机重启、网络抖动补拉这些情况。另外一个容易被忽视的点是消费顺序。Kafka单分区内消息有序但如果消费者是多线程处理乱序就会发生。像“创建订单”和“取消订单”两条消息如果取消先被执行创建后被执行数据就错了。这时通常的做法是把同一条业务流水发到同一个分区消费端用单线程或者按业务ID哈希后再串行处理。除非业务能容忍乱序否则不要为了吞吐量牺牲顺序。6. 进阶无锁队列和单调队列面试里的加分项6.1 无锁队列与CAS热搜词里有“c原子操作与无锁队列”说明不少同学也在关注无锁队列。Java里对应的是ConcurrentLinkedQueue它不依赖Synchronized或Lock而是利用AtomicReference的CAS操作实现并发安全。无锁的核心是乐观锁思想每次修改前先看看内存里期望的值是否还是原来的值如果是就CAS更新成功如果不是说明被别的线程改了那就自旋重试。// 伪代码offer过程的核心逻辑 public boolean offer(E e) { NodeE newNode new Node(e); for (;;) { NodeE t tail.get(); NodeE next t.next.get(); if (t tail.get()) { if (next null) { if (t.next.compareAndSet(null, newNode)) { tail.compareAndSet(t, newNode); return true; } } else { tail.compareAndSet(t, next); } } } }无锁队列的优点是避免了锁竞争带来的线程阻塞和上下文切换在高并发、短临界区的场景下有优势。但它不是银弹CAS自旋如果碰撞率高CPU消耗反而更大实现复杂对内存模型要求高自己手写很容易写出ABA问题ConcurrentLinkedQueue不是阻塞队列不能实现“队列满时生产者等待”这样的语义。所以我的建议是业务代码优先用BlockingQueue只有压测证明锁竞争是瓶颈才去考虑无锁优化。面试时能把CAS和ABA问题讲清楚已经比很多人强了。6.2 单调队列的实战形状前面在双端队列里提过滑动窗口这里再补一个更复杂的形状——单调队列不只是求最大值还可以维护“局部有序性”。比如有一系列入账记录窗口长度为5需要实时知道当前窗口内的最小值和最大值。用两个单调队列可以同时维护。关键是队列里存的是下标而不是值这样能精确判断元素是否还在窗口内。还有一个很典型的应用是“最长连续递增子序列”的变体需要维护一个窗口内的最值并配合双指针移动。如果你已经把ArrayDeque的pollFirst、pollLast、offerLast练熟这类题基本可以手写。在实现单调队列时有几个容易踩的坑出队的下标判断要写在最前面因为如果队首已经滑出窗口要先清除再考虑维护单调性。比较时注意等于号。求最大值时如果新元素和队尾元素相等通常要把队尾元素弹出因为更低的下标会先滑出窗口保留新下标更稳妥。初始化窗口时不要直接套用完整窗口的逻辑。我习惯先单独填充前k个元素然后再进入滑动循环代码结构更清晰不容易数组越界。ArrayDequeInteger q new ArrayDeque(); // 先用前k个元素初始化单调递减队列 for (int i 0; i k; i) { while (!q.isEmpty() nums[q.peekLast()] nums[i]) { q.pollLast(); } q.offerLast(i); } // 再滑动 for (int i k; i n; i) { // 取结果 while (!q.isEmpty() q.peekFirst() i - k 1) q.pollFirst(); while (!q.isEmpty() nums[q.peekLast()] nums[i]) q.pollLast(); q.offerLast(i); }这套模板我前后手写过不下十遍面试时从来没有在这个问题上卡过。6.3 我个人在实际项目里的选用习惯最后分享一点我自己的实践体会。很多人看队列的源码时容易陷进各种并发细节出不来比如LinkedBlockingQueue为什么用两把锁SynchronousQueue的栈和队列两种模式有什么区别。我觉得更好的学习路径是先搞懂业务场景再回头读源码。我的选用习惯可以概括成一句话能用有界就不要用无界能用阻塞就不要用空转能靠消息中间件解耦就不要在JVM里硬扛。比如内部异步任务我首选ArrayBlockingQueue并明确容量需要跨实例解耦我会直接引入MQ而不是在应用里手动搞同步双写在线程池配置文件里我会把线程池的队列类型、容量、拒绝策略都写进配置中心方便动态调整而不是在代码里硬编码。这样即使某天流量翻倍也能在不发版的情况下把队列容量调大。队列这东西看起来简单但几乎每一个线上故障里都能找到它的影子。希望大家能从我的经历里得到一点启发下次写Executors.newFixedThreadPool的时候停顿一下想想你用的那个队列是不是无界的下次用BlockingQueue.take()的时候想一想中断之后该怎么处理。这些细节才是Java开发真正拉开差距的地方。
返回列表