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

如何实现ThreadPoolExecutor按规则出队?同属性任务串行且限总线程数

实现同属性任务串行执行+总线程数控制的方案

核心需求回顾

  • 同color属性的任务必须按接收顺序串行执行
  • 不同color任务可并行执行,但总运行线程数有上限
  • 优先执行最旧的、无同color任务正在运行的任务

方案一:自定义阻塞队列+ThreadPoolExecutor

通过重写阻塞队列的出队逻辑,结合线程安全的运行状态跟踪,实现符合要求的任务调度。

1. 定义带color属性的任务类

封装原始任务,携带color标识,并在任务完成后释放该color的运行权限:

class ColorTask implements Runnable {
    private final String color;
    private final Runnable delegate;

    public ColorTask(String color, Runnable delegate) {
        this.color = color;
        this.delegate = delegate;
    }

    public String getColor() {
        return color;
    }

    @Override
    public void run() {
        try {
            delegate.run();
        } finally {
            RunningColors.release(color);
        }
    }
}

2. 跟踪正在运行的color状态

用线程安全的集合维护当前正在执行任务的color,确保同color任务不会并行:

class RunningColors {
    private static final ConcurrentHashMap<String, Integer> RUNNING_COLORS = new ConcurrentHashMap<>();

    // 尝试获取color的执行权限:若该color未在运行,则标记为运行中
    public static boolean tryAcquire(String color) {
        return RUNNING_COLORS.compute(color, (k, existing) -> existing == null ? 1 : null) != null;
    }

    // 释放color的执行权限
    public static void release(String color) {
        RUNNING_COLORS.remove(color);
    }
}

3. 自定义阻塞队列重写出队逻辑

重写take()方法,遍历队列找到第一个符合条件的任务(无同color任务在运行),同时保证队列其他元素顺序不变:

class ColorAwareBlockingQueue extends LinkedBlockingQueue<Runnable> {
    @Override
    public Runnable take() throws InterruptedException {
        while (true) {
            synchronized (this) {
                Iterator<Runnable> iterator = iterator();
                while (iterator.hasNext()) {
                    Runnable next = iterator.next();
                    if (next instanceof ColorTask) {
                        ColorTask task = (ColorTask) next;
                        if (RunningColors.tryAcquire(task.getColor())) {
                            iterator.remove();
                            return task;
                        }
                    } else {
                        // 非ColorTask直接取出执行
                        iterator.remove();
                        return next;
                    }
                }
            }
            // 队列中无符合条件的任务,等待新任务加入后重试
            wait();
        }
    }

    @Override
    public boolean offer(Runnable e) {
        boolean result = super.offer(e);
        notify(); // 有新任务加入,唤醒等待的出队线程
        return result;
    }
}

4. 创建自定义线程池

使用自定义队列初始化ThreadPoolExecutor,设置总线程数上限:

ThreadPoolExecutor executor = new ThreadPoolExecutor(
        2, // 核心线程数
        5, // 最大线程数(总运行线程上限)
        60L, TimeUnit.SECONDS,
        new ColorAwareBlockingQueue()
);

方案二:单线程池+Semaphore(更优方案)

利用单线程池天然保证同color任务串行,结合Semaphore控制总并发线程数,逻辑更简洁,维护成本更低。

实现代码

class ColorTaskManager {
    // 为每个color维护一个单线程池,确保同color任务串行
    private final ConcurrentHashMap<String, ExecutorService> colorExecutors = new ConcurrentHashMap<>();
    // 控制总运行线程数的信号量
    private final Semaphore globalThreadLimit;

    public ColorTaskManager(int maxTotalThreads) {
        this.globalThreadLimit = new Semaphore(maxTotalThreads);
    }

    public void submitTask(String color, Runnable task) {
        // 懒加载对应color的单线程池
        ExecutorService colorExecutor = colorExecutors.computeIfAbsent(
                color, k -> Executors.newSingleThreadExecutor()
        );

        colorExecutor.submit(() -> {
            try {
                // 获取全局线程权限,达到上限则阻塞
                globalThreadLimit.acquire();
                task.run();
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            } finally {
                // 释放全局线程权限
                globalThreadLimit.release();
            }
        });
    }

    // 关闭所有线程池
    public void shutdownAll() {
        colorExecutors.values().forEach(ExecutorService::shutdown);
    }
}

方案优势

  • 无需自定义队列和线程池,复用JDK原生组件,稳定性更高
  • 逻辑清晰:单线程池保证同color串行,Semaphore控制总并发,避免复杂的队列遍历逻辑
  • 易于扩展:可针对不同color添加额外的调度规则

方案对比

方案类型优点缺点
自定义队列+ThreadPoolExecutor贴近原生线程池模型,适合对线程池行为有强定制需求的场景队列遍历和同步逻辑复杂,高并发下可能有性能损耗
单线程池+Semaphore逻辑简洁,维护成本低,性能稳定线程池数量随color数量增长,但单线程池资源开销极小,可忽略

内容的提问来源于stack exchange,提问作者Jack Cole

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 23:55:39