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

基于公平性的多客户多任务流式调度器实现咨询

调度器多线程实现方案指导

核心需求明确

  • 流式任务队列:任务持续动态流入,而非预定义集合
  • 客户任务串行:同一客户的任务同一时间仅能执行一个,避免同一客户任务并发
  • 跨客户并行:不同客户的任务可同时执行,最大化利用线程资源

关键实现思路

  1. 线程池管理并发:用线程池统一管理执行线程,控制全局并发数,避免手动创建线程的资源浪费和风险
  2. 客户任务队列隔离:为每个客户维护独立的任务队列,确保同一客户的任务按顺序排队
  3. 执行状态控制:标记每个客户是否有任务正在执行,避免重复提交同一客户的任务,同时在任务执行完毕后自动触发队列中剩余任务的执行

具体代码实现示例(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 01:10:18