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
相关产品推荐
相关产品推荐

