You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 06:55:16