本文介绍一种高效、可扩展的并发调度方案:为相同属性(如“color”)的事件分配独立队列并由专用工作线程串行处理,同时通过共享线程池实现全局并发数限制,兼顾顺序性、隔离性与资源利用率。
本文介绍一种高效、可扩展的并发调度方案:为相同属性(如“color”)的事件分配独立队列并由专用工作线程串行处理,同时通过共享线程池实现全局并发数限制,兼顾顺序性、隔离性与资源利用率。
在高吞吐事件处理场景中,常需满足“同组串行、跨组并发、全局限流”三重要求——例如按 color 字段分组,绿色事件必须严格按接收顺序依次执行,黄色事件亦然;但绿与黄之间无需顺序约束;且所有颜色的执行线程总数不能超过系统设定上限(如 8 个并发线程)。
直接改造 ThreadPoolExecutor 的任务队列(如自定义 BlockingQueue)来实现“跳过同色正在运行任务”的动态出队逻辑,不仅破坏线程池设计契约,还极易引发竞态、死锁或饥饿问题(如某色任务持续积压导致其他色长期得不到调度)。因此,推荐采用“分组队列 + 共享工作者池”的解耦架构:
// 1. 分组队列容器private final ConcurrentMap<String, Queue<Event>> colorQueues = new ConcurrentHashMap<>();private final ReentrantLock lock = new ReentrantLock();private final Set<String> runningColors = ConcurrentHashMap.newKeySet();// 2. 工作线程任务(提交至共享线程池)Runnable workerTask = () -> { while (!Thread.currentThread().isInterrupted()) { Event event = null; String color = null; // 轮询所有队列,找到首个可执行的 oldest 事件(避免饿死) for (Queue<Event> queue : colorQueues.values()) { if (!queue.isEmpty()) { event = queue.peek(); // 不移除,先检查 if (event != null && !runningColors.contains(event.color())) { color = event.color(); event = queue.poll(); // 确认后出队 break; } } } if (event == null) { Thread.sleep(10); // 短暂让出 CPU continue; } // 标记 color 正在运行 runningColors.add(color); try { event.execute(); // 执行业务逻辑 } finally { runningColors.remove(color); // 必须确保释放 } }};// 启动 N 个 worker 线程ExecutorService workers = Executors.newFixedThreadPool(8);for (int i = 0; i < 8; i++) { workers.submit(workerTask);}
该方案天然满足所有原始需求:同色严格 FIFO、跨色完全并发、全局线程数可控,且代码清晰、易于监控与调试。相比侵入式修改线程池队列,它更符合面向对象与关注点分离原则,是生产环境推荐的稳健实践。