基于公平性的多客户多任务流式调度器实现咨询
调度器多线程实现方案指导
核心需求明确
- 流式任务队列:任务持续动态流入,而非预定义集合
- 客户任务串行:同一客户的任务同一时间仅能执行一个,避免同一客户任务并发
- 跨客户并行:不同客户的任务可同时执行,最大化利用线程资源
关键实现思路
- 线程池管理并发:用线程池统一管理执行线程,控制全局并发数,避免手动创建线程的资源浪费和风险
- 客户任务队列隔离:为每个客户维护独立的任务队列,确保同一客户的任务按顺序排队
- 执行状态控制:标记每个客户是否有任务正在执行,避免重复提交同一客户的任务,同时在任务执行完毕后自动触发队列中剩余任务的执行
具体代码实现示例(Java)
任务类定义
class MyJob implements Runnable { private final Integer customerId; private final int taskParam; public MyJob(Integer customerId, int taskParam) { this.customerId = customerId; this.taskParam = taskParam; } @Override public void run() { // 替换为实际任务执行逻辑 System.out.printf("客户%d的任务执行中,参数:%d%n", customerId, taskParam); try { Thread.sleep(100); // 模拟任务耗时 } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }
调度器核心实现
import java.util.Queue; import java.util.concurrent.*; class CustomerJobScheduler { // 根据系统资源调整线程池大小,这里用固定大小线程池示例 private final ExecutorService executor = Executors.newFixedThreadPool(10); // 线程安全的映射:客户ID -> 该客户的待执行任务队列 private final ConcurrentHashMap<Integer, Queue<Runnable>> customerJobQueues = new ConcurrentHashMap<>(); // 标记客户是否有任务正在执行,避免重复触发 private final ConcurrentHashMap<Integer, Boolean> customerRunningStatus = new ConcurrentHashMap<>(); // 外部调用此方法提交流式任务 public void submitJob(Integer customerId, Runnable job) { // 获取或创建客户的任务队列 Queue<Runnable> jobQueue = customerJobQueues.computeIfAbsent(customerId, k -> new ConcurrentLinkedQueue<>()); jobQueue.add(job); // 尝试触发该客户的任务执行 triggerTaskExecution(customerId); } private void triggerTaskExecution(Integer customerId) { // 原子操作:如果客户未在执行任务,则标记为正在执行并提交任务 if (customerRunningStatus.putIfAbsent(customerId, Boolean.TRUE) == null) { executor.submit(() -> { try { Queue<Runnable> jobQueue = customerJobQueues.get(customerId); if (jobQueue == null) return; // 循环执行队列中的任务,直到为空 Runnable currentJob; while ((currentJob = jobQueue.poll()) != null) { currentJob.run(); } } finally { // 执行完毕后移除运行标记 customerRunningStatus.remove(customerId); // 再次检查队列,防止执行过程中有新任务流入 Queue<Runnable> jobQueue = customerJobQueues.get(customerId); if (jobQueue != null && !jobQueue.isEmpty()) { triggerTaskExecution(customerId); } } }); } } // 优雅关闭调度器 public void shutdown() { executor.shutdown(); try { if (!executor.awaitTermination(60, TimeUnit.SECONDS)) { executor.shutdownNow(); } } catch (InterruptedException e) { executor.shutdownNow(); } } }
使用示例
public class SchedulerDemo { public static void main(String[] args) throws InterruptedException { CustomerJobScheduler scheduler = new CustomerJobScheduler(); // 模拟流式任务提交:分批、动态添加 scheduler.submitJob(1, new MyJob(1, 100)); scheduler.submitJob(1, new MyJob(1, 200)); scheduler.submitJob(2, new MyJob(2, 300)); Thread.sleep(200); // 模拟任务流入间隔 scheduler.submitJob(1, new MyJob(1, 400)); scheduler.submitJob(2, new MyJob(2, 500)); // 等待所有任务执行完成 Thread.sleep(1000); scheduler.shutdown(); } }
注意事项
- 线程池参数调优:根据服务器CPU核心数、任务耗时调整线程池大小,可使用
ThreadPoolExecutor自定义核心线程数、最大线程数、任务队列等参数 - 内存监控:若某客户任务流入速度远快于执行速度,会导致其任务队列持续膨胀,需添加限流或告警机制
- 异常处理:任务执行时需捕获并处理异常,避免单个任务失败导致后续任务无法执行
- 优雅停机:关闭调度器时需等待已提交任务执行完成,或根据业务需求选择是否中断未执行任务
内容的提问来源于stack exchange,提问作者Himanshu Goyal
相关产品推荐
相关产品推荐

