ThreadPoolExecutor队列满时优先执行队列任务而非新提交任务的实现方案
ThreadPoolExecutor任务执行优先级问题
问题描述
使用有界队列的ThreadPoolExecutor默认执行逻辑如下:
- 核心线程(
corePoolSize)繁忙时,新任务加入队列 - 队列已满且最大线程数(
maximumPoolSize)大于核心线程数时,新任务会触发线程扩容,直接由新建线程执行,而非等待队列中的先提交任务完成
这种逻辑会导致后提交的任务先于队列中先提交的任务执行,极端情况下若任务提交速度与max线程执行速度持平,队列中的任务可能永远无法被执行,不符合“先提交先执行”的预期。
测试验证
测试代码
package com.multithreading; import java.util.concurrent.*; import java.util.function.Function; public class ExecutorOrder { public static void main(String[] args) throws InterruptedException { BlockingQueue<Runnable> blockingQueue = new ArrayBlockingQueue<>(2); ExecutorService e = new ThreadPoolExecutor(2, 4, 5, TimeUnit.SECONDS, blockingQueue); Function<Integer, Runnable> threadGenerator = (Integer t) -> () -> { System.out.println("This is task " + t + ", now sleeping for 2 seconds: " + Thread.currentThread().getName()); try { Thread.sleep(2000); } catch (InterruptedException ex) { throw new RuntimeException(ex); } System.out.println("Task " + t + " exits, " + Thread.currentThread().getName()); }; e.submit(threadGenerator.apply(1)); e.submit(threadGenerator.apply(2)); e.submit(threadGenerator.apply(3)); e.submit(threadGenerator.apply(4)); e.submit(threadGenerator.apply(5)); e.submit(threadGenerator.apply(6)); Thread.sleep(3000); e.shutdown(); } }
执行流程
- 任务1、2占用核心线程开始执行
- 任务3、4提交时核心线程繁忙,进入队列,队列达到容量上限
- 任务5、6提交时队列已满,触发线程扩容至
maximumPoolSize,直接由新建线程执行,队列中的3、4延后处理
示例输出
This is task 1, now sleeping for 2 seconds: pool-1-thread-1 This is task 6, now sleeping for 2 seconds: pool-1-thread-4 This is task 5, now sleeping for 2 seconds: pool-1-thread-3 This is task 2, now sleeping for 2 seconds: pool-1-thread-2 Task 5 exits, pool-1-thread-3 Task 1 exits, pool-1-thread-1 Task 2 exits, pool-1-thread-2 Task 6 exits, pool-1-thread-4 This is task 3, now sleeping for 2 seconds: pool-1-thread-2 This is task 4, now sleeping for 2 seconds: pool-1-thread-4 Task 4 exits, pool-1-thread-4 Task 3 exits, pool-1-thread-2
解决方案
要实现优先执行队列中先提交的任务,新任务始终入队等待,可以通过以下两种方式实现:
方式一:重写ThreadPoolExecutor的execute方法
修改队列满时的处理逻辑,阻塞等待队列空闲后再入队,而非新建线程执行新任务:
public class PriorityQueueThreadPoolExecutor extends ThreadPoolExecutor { public PriorityQueueThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueue<Runnable> workQueue) { super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue); } @Override public void execute(Runnable command) { if (command == null) { throw new NullPointerException(); } int c = ctl.get(); // 核心线程未满,直接新建线程执行 if (workerCountOf(c) < getCorePoolSize()) { if (addWorker(command, true)) { return; } c = ctl.get(); } // 核心线程已满,尝试入队 if (isRunning(c) && getQueue().offer(command)) { int recheck = ctl.get(); // 线程池已停止,移除任务并拒绝 if (!isRunning(recheck) && remove(command)) { reject(command); } // 无存活线程,新建线程处理队列任务 else if (workerCountOf(recheck) == 0) { addWorker(null, false); } } else { // 队列满时,阻塞等待队列有空间再入队 try { getQueue().put(command); } catch (InterruptedException e) { Thread.currentThread().interrupt(); reject(command); } } } }
方式二:使用带优先级的队列
为任务添加提交顺序作为优先级,确保先提交的任务被优先取出执行:
// 包装任务,记录提交顺序 static class OrderedTask implements Runnable, Comparable<OrderedTask> { private final Runnable task; private final long order; private static final AtomicLong sequence = new AtomicLong(0); public OrderedTask(Runnable task) { this.task = task; this.order = sequence.getAndIncrement(); } @Override public void run() { task.run(); } @Override public int compareTo(OrderedTask o) { return Long.compare(this.order, o.order); } } // 使用PriorityBlockingQueue作为任务队列 BlockingQueue<Runnable> queue = new PriorityBlockingQueue<>(); ExecutorService e = new ThreadPoolExecutor(2, 4, 5, TimeUnit.SECONDS, queue); // 提交任务时包装为OrderedTask e.submit(new OrderedTask(threadGenerator.apply(1))); // 其他任务同理包装后提交
内容的提问来源于stack exchange,提问作者Vaibhav
相关产品推荐
相关产品推荐

