如何避免ExecutorService中重复任务引发的线程阻塞?
问题
我有一个会多次执行且传入不同字符串参数的任务,希望避免相同字符串参数的两个任务同时执行。通过ExecutorService提交这些任务时,目前出现了这样的问题:第二个"hi"任务会阻塞其他任务,导致仅4/5个任务并发执行;待第一个"hi"任务释放锁后,第五个任务才继续执行,其余任务后续运行正常。想请教是否有办法避免此类任务阻塞,让其他3个任务优先执行,直到真正有5个任务并发时再排队?
任务提交代码:
executor.submit(new Task("hi")); executor.submit(new Task("h")); executor.submit(new Task("u")); executor.submit(new Task("y")); executor.submit(new Task("hi")); executor.submit(new Task("p")); executor.submit(new Task("o")); executor.submit(new Task("bb"));
任务实现代码:
Lock l = getLock(x); try { l.lock(); System.out.println(x); try { Thread.sleep(5000); } catch (InterruptedException ex) { Logger.getLogger(Task.class.getName()).log(Level.SEVERE, null, ex); } } finally { l.unlock(); }
解决方案
核心问题是当前锁机制会让相同参数的任务直接占用线程池线程并阻塞,浪费核心线程资源。要解决这个问题,需让相同参数的等待任务不占用线程池线程,而是进入参数专属的等待队列,待该参数任务执行完成后再取出队列任务执行。
可以通过ConcurrentHashMap结合Semaphore实现这个逻辑:
实现思路
- 给每个字符串参数分配一个许可数为1的Semaphore,保证同一参数任务互斥执行
- 用队列缓存相同参数的等待任务,避免占用线程池线程
- 当前参数任务执行完毕后,自动触发队列中等待任务的执行
代码示例
任务调度器类
import java.util.concurrent.*; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; public class TaskScheduler { private final ExecutorService executor; private final ConcurrentHashMap<String, Semaphore> semaphoreMap = new ConcurrentHashMap<>(); private final ConcurrentHashMap<String, Queue<Runnable>> taskQueueMap = new ConcurrentHashMap<>(); private final Lock queueLock = new ReentrantLock(); public TaskScheduler(ExecutorService executor) { this.executor = executor; } public void submitTask(String param, Runnable task) { // 获取或创建当前参数对应的Semaphore Semaphore semaphore = semaphoreMap.computeIfAbsent(param, k -> new Semaphore(1)); if (semaphore.tryAcquire()) { // 获取许可成功,直接提交任务执行 executor.submit(() -> { try { task.run(); } finally { semaphore.release(); // 释放许可后处理该参数的等待任务 processWaitingTasks(param); } }); } else { // 未获取到许可,将任务加入等待队列 queueLock.lock(); try { taskQueueMap.computeIfAbsent(param, k -> new LinkedList<>()).offer(task); } finally { queueLock.unlock(); } } } private void processWaitingTasks(String param) { queueLock.lock(); try { Queue<Runnable> queue = taskQueueMap.get(param); if (queue != null && !queue.isEmpty()) { Runnable nextTask = queue.poll(); // 递归提交队列中的任务 submitTask(param, nextTask); } else { // 队列空了,清理缓存的Semaphore和队列(可选,避免内存占用) taskQueueMap.remove(param); semaphoreMap.remove(param); } } finally { queueLock.unlock(); } } }
使用方式
// 初始化核心线程数为5的线程池 ExecutorService executor = Executors.newFixedThreadPool(5); TaskScheduler scheduler = new TaskScheduler(executor); // 提交任务 scheduler.submitTask("hi", () -> { System.out.println("hi"); try { Thread.sleep(5000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }); scheduler.submitTask("h", () -> { System.out.println("h"); try { Thread.sleep(5000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }); // 其余任务按同样方式提交...
效果说明
- 线程池的5个核心线程会优先处理不同参数的任务,不会因为相同参数的等待任务占用线程
- 只有当所有核心线程都被不同参数的任务占用时,新任务才会进入线程池的等待队列
- 相同参数的任务会按提交顺序依次执行,不会出现并发冲突
内容的提问来源于stack exchange,提问作者Jetski111
相关产品推荐
相关产品推荐

