如何实现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
相关产品推荐
相关产品推荐

