【问题标题】:Spring Boot:How can we implement multiple @Scheduled tasks with each having its own thread pool?Spring Boot:我们如何实现多个@Scheduled 任务,每个任务都有自己的线程池?
【发布时间】:2020-03-07 22:03:53
【问题描述】:

我想实现多个@Scheduled(具有固定延迟)任务,每个任务都有自己的线程池。

@Scheduled(fixedDelayString = "30000")
public void createOrderSchedule() {
    //create 10 orders concurrently; wait for all to be finished
    createOrder(10);
}

@Scheduled(fixedDelayString = "30000")
public void processOrderSchedule() {
    //process 10 orders concurrently; wait for all to be finished
}

@Scheduled(fixedDelayString = "30000")
public void notifySchedule() {
    //send notification for 10 orders concurrently; wait for all to be finished
}

我设法为每个调度程序创建了不同的ThreadPoolTaskExecutor,如下所示:

@Bean("orderPool")
public ThreadPoolTaskExecutor createOrderTaskExecutor() {
    ThreadPoolTaskExecutor pool = new ThreadPoolTaskExecutor();
    pool.setCorePoolSize(5);
    pool.setMaxPoolSize(10);
    pool.setThreadNamePrefix("order-thread-pool-");
    pool.setWaitForTasksToCompleteOnShutdown(true);
    return pool;
}

..

我为每个任务提供了@Async

@Async("orderPool")
public void createOrder(Integer noOforders) {..}

和一个任务调度器配置

@Bean
public ThreadPoolTaskScheduler threadPoolTaskScheduler() {
    ThreadPoolTaskScheduler threadPoolTaskScheduler = new ThreadPoolTaskScheduler();
    threadPoolTaskScheduler.setPoolSize(3);
    return threadPoolTaskScheduler;
}

我使用CompletableFuture.allOf(..).join(); 等待每个任务完成,但它阻止了所有其他@Scheduled 任务。

综上所述,我想实现以下目标:

  1. 每个@Scheduled 任务都应独立运行,而不会阻塞其他@Scheduled 任务。
  2. 每个@Scheduled 任务都应该有自己的线程池,以便它可以同时运行多个子任务(比如10 个)。
  3. 每个@Scheduled 任务都必须等待每个触发器完成而不被再次调用。

我怎样才能做到这一点?

【问题讨论】:

标签: multithreading spring-boot scheduled-tasks threadpool spring-scheduled


【解决方案1】:

在连续使用了将近 18 个小时之后,我能够实现我在​​上述问题中提出的问题。抱歉这么晚了。

所以流 API 提供了像 IntStream 等接口来并行地流式传输元素。这导致我并行创建n 订单。 (同时在不同的调度器中并行处理k的订单。以此类推。)

IntStream.range(0, inputIds.size())
        .parallel().forEach(index -> createOrder(inputIds.get(index)));

就这么简单。解决了 1 个用例。现在我希望这个调度程序有它自己的池。发现IntStream.parallel()使用ForkJoinPool,它是我们自己的ExecutorService的继承者,令我惊讶的是,spring为ForkJoinPool提供了一个预配置的工厂bean,即ForkJoinPoolFactoryBean。所以我创建了一个名为createOrderExecutor的bean。

@Bean("createOrderExecutor")
public ForkJoinPoolFactoryBean createOrderExecutor() {
    ForkJoinPoolFactoryBean createOrderPoolFactoryBean = new ForkJoinPoolFactoryBean();
    createOrderPoolFactoryBean.setParallelism(10);
    createOrderPoolFactoryBean.setAsyncMode(true);
    createOrderPoolFactoryBean.setUncaughtExceptionHandler(null);
    createOrderPoolFactoryBean.setThreadFactory(p -> {
        final ForkJoinWorkerThread worker = ForkJoinPool.defaultForkJoinWorkerThreadFactory.newThread(p);
        worker.setName("create-order-pool-" + worker.getPoolIndex());
        return worker;
    });
    return createOrderPoolFactoryBean;
}

我在调度程序类中自动装配了这个 bean,并同时提交了所有订单,如下所示。

createOrderExecutor.getObject().submit(() -> IntStream.range(0, inputIds.size())
                .parallel().forEach(index -> createOrder(inputIds.get(index))));

那里。解决了第二个用例。现在这不会等待所有并行任务完成,它只会异步触发它们。现在,ForkJoinTask(即submit() 返回)提供了一个get() 方法,该方法等待计算完成并返回结果。 (但我不需要结果,我宁愿用try-catch包围它们。而且,我会等待完成。)

@Scheduled(fixedDelayString = "5000")
public void createOrderScheduler() {
        createOrderExecutor.getObject().submit(() -> IntStream.range(0, inputIds.size())
                .parallel().forEach(index -> createOrder(inputIds.get(index)))).get();
}

这解决了我的最后一个用例。我为应用程序中的所有调度程序都这样做了。

相信我,我尝试了几乎所有在线提供的 CompletableFuture 实现,但都无法实现所有这些。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2023-04-04
    • 2023-01-30
    • 1970-01-01
    • 2012-09-03
    • 2016-08-21
    • 2015-04-24
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多