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

Spring Boot中如何让下一次@Scheduler调用等待当前任务完成?

Spring Boot调度任务与异步处理重复记录问题解决方案

问题场景

我正在开发Spring Boot应用,使用@Async多线程异步处理大量记录:将主列表拆分为8个子列表,交由8个线程并行处理。同时通过@Scheduler配置每2秒触发一次任务。

但存在核心问题:由于调度间隔较短,当前任务尚未完成时,下一次调度就会触发,导致记录重复处理。例如首次调度从数据库查询flag为0的72000条记录,异步任务处理时会将这些记录的flag改为1,但2秒后下一次调度触发时,部分正在处理的记录还未完成flag更新,会被再次查询出来,造成重复处理。从日志可见,首次调度获取72000条记录后,异步线程开始处理,此时下一次调度触发并获取16000条记录,其中包含正在处理的记录。

核心需求

下一次@Scheduler调用必须等待当前调度的所有异步任务完成后再执行,且无法增加调度间隔——因为数据量波动较大,有时仅400-500条,有时达数千条。

解决方案

关键修改点

  • 调整异步方法的返回值为Future<Void>,用于跟踪异步任务的执行状态
  • 在调度方法中收集所有异步任务的Future对象,等待全部任务完成后再结束当前调度方法(Spring默认调度器为单线程,当前调度方法未结束时,下一次调度会被阻塞)
  • 优化线程池配置(可选),避免任务队列溢出导致拒绝执行

修改后的代码

1. 主启动类(DemoApplication)

@SpringBootApplication
@EnableScheduling
@EnableAsync
public class DemoApplication {

    public static void main(String[] args) {
        SpringApplication.run(DemoApplication.class, args);
    }
    
    @Bean
    public RestTemplate restTemplate(RestTemplateBuilder builder) {
        return builder.build();
    }
    
    @Bean("ThreadPoolTaskExecutor")
    public TaskExecutor getAsyncExecutor() {
        final ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(8);
        executor.setMaxPoolSize(8);
        executor.setQueueCapacity(100); // 适当调大队列容量,避免任务被拒绝
        executor.setWaitForTasksToCompleteOnShutdown(true);
        executor.setThreadNamePrefix("async-");
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); // 改用CallerRunsPolicy,避免任务丢失
        executor.initialize();
        return executor;
    }
}

2. 业务服务类(Main)

@Service
public class Main {

    @Autowired
    private DemoDao dao;
    
    static int schedulerCount = 0;
    
    @Scheduled(cron = "0/2 * * * * *")
    public void schedule() {
        System.out.println("++++++++++++++++++++++++++++++++++++Scheduler started schedulerCount  :  "+schedulerCount+"+++++++++++++++++++++++++++++++++++++"+ LocalDateTime.now());
        List<Json> jsonList = new ArrayList<>();
        List<List<Json>> smallLi = new ArrayList<>();
        // 用于收集所有异步任务的Future
        List<Future<Void>> futures = new ArrayList<>();
        try {
            jsonList = dao.getJsonList();
            System.out.println("jsonList size   :   " + jsonList.size());
            
            int count = jsonList.size();
            if (count == 0) {
                schedulerCount++;
                return; // 无数据时直接返回,避免空处理
            }
            
            // 拆分主列表为8个子列表,用整数除法避免浮点运算误差
            int limit = (count + 7) / 8;
            for (int j = 0; j < count; j += limit) {
                smallLi.add(new ArrayList<>(jsonList.subList(j, Math.min(count, j + limit))));
            }
            System.out.println("smallLi  :  " + smallLi.size());
            
            // 提交异步任务并收集Future
            for (List<Json> subList : smallLi) {
                futures.add(withAsyn(subList, schedulerCount));
            }
            
            // 等待所有异步任务完成
            for (Future<Void> future : futures) {
                try {
                    future.get();
                } catch (InterruptedException | ExecutionException e) {
                    e.printStackTrace();
                }
            }
            
            schedulerCount++;
        } catch (Exception e) {
            e.printStackTrace();
        }
    }

    @Async("ThreadPoolTaskExecutor")
    public Future<Void> withAsyn(List<Json> li, int schedulerCount) throws Exception {
        System.out.println("with start+++++++++++++ schedulerCount " + schedulerCount + ", name  :  "
                + Thread.currentThread().getName() + ", time  :  " + LocalDateTime.now() + ", start index  :  "
                + li.get(0).getId() + ", end index  :  " + li.get(li.size() - 1).getId());
        try {
            XSSFWorkbook workbook = new XSSFWorkbook();
            XSSFSheet spreadsheet = workbook.createSheet("Data");
            XSSFRow row;

            for (int i = 0; i < li.size(); i++) {
                row = spreadsheet.createRow(i);
                Cell cell9 = row.createCell(0);
                cell9.setCellValue(li.get(i).getId());

                Cell cell = row.createCell(1);
                cell.setCellValue(li.get(i).getName());

                Cell cell1 = row.createCell(2);
                cell1.setCellValue(li.get(i).getPhone());

                Cell cell2 = row.createCell(3);
                cell2.setCellValue(li.get(i).getEmail());

                Cell cell3 = row.createCell(4);
                cell3.setCellValue(li.get(i).getAddress());

                Cell cell4 = row.createCell(5);
                cell4.setCellValue(li.get(i).getPostalZip());
            }
            FileOutputStream out = new FileOutputStream(new File("C:\\Users\\RK658\\Desktop\\logs\\generated\\"
                    + Thread.currentThread().getName() + "_" + schedulerCount + ".xlsx"));
            workbook.write(out);
            out.close();
            // 建议将flag更新逻辑放在这里,确保文件生成完成后再修改数据库状态
            // dao.updateFlagForRecords(li);
        } catch (Exception e) {
            e.printStackTrace();
        }
        System.out.println("with end+++++++++++++ schedulerCount " + schedulerCount + ", name  :  "
                + Thread.currentThread().getName() + ", time  :  " + LocalDateTime.now() + ", start index  :  "
                + li.get(0).getId() + ", end index  :  " + li.get(li.size() - 1).getId());
        return new AsyncResult<>(null);
    }
}

额外说明

  • 线程池配置中,将队列容量从8调大到100,并改用CallerRunsPolicy作为拒绝策略,避免当任务量突增时任务被直接拒绝。
  • 在调度方法中添加了空列表判断,无数据时直接返回,减少不必要的处理。
  • 建议将flag更新逻辑放在异步任务的最后,确保文件生成完成后再修改数据库状态,进一步降低重复处理的风险。

内容的提问来源于stack exchange,提问作者Rahul

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 04:42:55