Java自定义ThreadPoolExecutor中PriorityBlockingQueue排序失效求助
问题描述
尝试实现一个持有PriorityBlockingQueue的自定义ThreadPoolExecutor,期望队列按「数值越小优先级越高」的规则对任务排序,但测试发现任务按提交顺序执行,PriorityBlockingQueue的大小始终为0,且优先级值为2的IO任务先于优先级值为1的计算任务执行,不符合预期。
问题原因分析
- 队列排序规则未生效:初始化
PriorityBlockingQueue时未传入自定义Comparator,且Task类未实现Comparable接口,队列无法识别任务的优先级逻辑。 - 任务包装导致优先级丢失:调用
super.submit(task)时,ThreadPoolExecutor会将Task包装为FutureTask实例,而FutureTask不具备优先级比较能力,队列实际存储的是FutureTask而非原始Task,自然无法按预期排序。 - maxPriority逻辑错误:当前代码取任务优先级数值的最大值,但需求是「数值越小优先级越高」,且第三个提交的
OTHER类型任务(优先级3)导致maxPriority被错误设置为3。 - 任务直接被核心线程执行:若提交的任务数量小于线程池核心线程数,任务会被直接分配线程执行,不会进入队列,这也是队列大小始终为0的原因。
修复方案
1. 让Task实现Comparable接口
修改Task类,实现Comparable<Task>接口,定义优先级排序逻辑:
public class Task implements Callable<Object>, Comparable<Task> { private Callable<?> callable; private TaskType priority; public <T> Task(Callable<T> callable, TaskType priority) { this.callable = callable; this.priority = priority; } public TaskType getPriority() { return priority; } public static <T> Task createTask(Callable<T> task, TaskType priority) { return new Task(task, priority); } @Override public Object call() throws Exception { return callable.call(); } @Override public int compareTo(Task other) { // 数值越小优先级越高,当前任务优先级数值更小时返回-1,排在队列前面 return Integer.compare(this.priority.getPriorityValue(), other.priority.getPriorityValue()); } }
2. 重写ThreadPoolExecutor的newTaskFor方法
自定义FutureTask保留优先级比较能力,确保队列能识别任务优先级:
public class CustomExecutor extends ThreadPoolExecutor { private int highestPriority = Integer.MAX_VALUE; // 记录当前最高优先级(数值最小的) private boolean isShutdown = false; public CustomExecutor() { super(Runtime.getRuntime().availableProcessors() / 2, Runtime.getRuntime().availableProcessors() - 1, 300, TimeUnit.MILLISECONDS, new PriorityBlockingQueue<>()); // 现在Task实现Comparable,队列自动排序 } // 重写newTaskFor,返回支持优先级比较的FutureTask @Override protected <T> RunnableFuture<T> newTaskFor(Callable<T> callable) { if (callable instanceof Task) { return new PriorityFutureTask<>((Task) callable); } return super.newTaskFor(callable); } private <T> Future<T> submitTask(Task task) { // 更新最高优先级:仅当当前任务优先级数值更小时更新 if (task.getPriority().getPriorityValue() < highestPriority) { highestPriority = task.getPriority().getPriorityValue(); } if (isShutdown) { System.out.println("Thread pool has already been shut down. Cannot submit task"); return null; } return super.submit(task); } public <T> Future<T> submit(Callable<T> callable, TaskType priority) { Task task = new Task(callable, priority); return submitTask(task); } public <T> Future<T> submit(Callable<T> callable) { Task task = new Task(callable, TaskType.OTHER); return submitTask(task); } public int getCurrentHighestPriority() { return highestPriority; } public void gracefullyTerminate() { isShutdown = true; super.shutdown(); } // 自定义PriorityFutureTask,基于Task的优先级排序 private static class PriorityFutureTask<T> extends FutureTask<T> implements Comparable<PriorityFutureTask<T>> { private final Task task; public PriorityFutureTask(Task task) { super(task); this.task = task; } @Override public int compareTo(PriorityFutureTask<T> other) { return task.compareTo(other.task); } } }
3. 调整测试逻辑(可选)
为了明显观察队列排序效果,可以增加任务提交数量,或者减小核心线程数,确保部分任务进入队列等待执行。
验证修复效果
修复后,优先级数值为1的COMPUTATIONAL任务会优先于数值为2的IO任务执行,队列会按优先级规则排序,currentHighestPriority会正确记录当前最小的优先级数值(测试场景下应为1)。
内容的提问来源于stack exchange,提问作者Achiya Ben Natan
相关产品推荐
相关产品推荐

