Spring每日调度任务:如何获取线程池任务的成功失败统计矩阵
如何在并发Job中整合Reader与Processor的统计信息
这个需求我之前做批量数据同步任务时碰到过,刚好可以给你几个适配并发场景的实现思路,核心是用线程安全的统计容器来统一收集Reader的读取统计和Processor的处理统计,最后在Job结束时汇总出整体的成功/失败数据。
方案一:自定义线程安全的统计持有类
这是最直接的方式,创建一个全局的统计类,用线程安全的计数器来维护各项指标,让Reader和Processor都引用它,实时更新统计数据。
1. 定义统计持有类
import java.util.concurrent.atomic.AtomicInteger; public class JobStatsHolder { // Reader相关统计 private final AtomicInteger totalReadAttempts = new AtomicInteger(0); private final AtomicInteger successfulReads = new AtomicInteger(0); private final AtomicInteger failedReads = new AtomicInteger(0); // Processor相关统计 private final AtomicInteger totalProcessAttempts = new AtomicInteger(0); private final AtomicInteger successfulProcesses = new AtomicInteger(0); private final AtomicInteger failedProcesses = new AtomicInteger(0); // Reader统计更新方法 public void incrementTotalReads() { totalReadAttempts.incrementAndGet(); } public void incrementSuccessfulReads() { successfulReads.incrementAndGet(); } public void incrementFailedReads() { failedReads.incrementAndGet(); } // Processor统计更新方法 public void incrementTotalProcesses() { totalProcessAttempts.incrementAndGet(); } public void incrementSuccessfulProcesses() { successfulProcesses.incrementAndGet(); } public void incrementFailedProcesses() { failedProcesses.incrementAndGet(); } // 汇总打印统计信息 public void printStats() { System.out.println("=== Job 执行统计 ==="); System.out.printf("Reader 总尝试读取: %d, 成功: %d, 失败: %d%n", totalReadAttempts.get(), successfulReads.get(), failedReads.get()); System.out.printf("Processor 总尝试处理: %d, 成功: %d, 失败: %d%n", totalProcessAttempts.get(), successfulProcesses.get(), failedProcesses.get()); } // 重置统计数据,避免多次执行累积 public void reset() { totalReadAttempts.set(0); successfulReads.set(0); failedReads.set(0); totalProcessAttempts.set(0); successfulProcesses.set(0); failedProcesses.set(0); } }
2. 修改Reader类,加入统计逻辑
假设你的Reader是分批读取API数据,每次读取一批后更新统计:
@Component public class Reader { @Autowired private JobStatsHolder statsHolder; public List<DataItem> readBatch(int batchSize) { statsHolder.incrementTotalReads(); try { // 调用API读取数据的逻辑 List<DataItem> dataItems = apiClient.fetchData(batchSize); statsHolder.incrementSuccessfulReads(); return dataItems; } catch (ApiException e) { statsHolder.incrementFailedReads(); // 可根据需求选择抛出异常或返回空列表 return Collections.emptyList(); } } }
3. 修改Processor任务,加入统计逻辑
每个线程池中的任务处理数据时更新统计:
@Component public class DataProcessor { @Autowired private JobStatsHolder statsHolder; @Autowired private DataRepository dataRepository; public void processData(DataItem item) { statsHolder.incrementTotalProcesses(); try { // 写入数据库的逻辑 dataRepository.save(item); statsHolder.incrementSuccessfulProcesses(); } catch (DbException e) { statsHolder.incrementFailedProcesses(); // 记录日志或做失败重试逻辑 log.error("处理数据失败: {}", item.getId(), e); } } }
4. 在JobHandler中初始化统计并汇总
@Component public class JobHandler { @Autowired private Reader reader; @Autowired private DataProcessor dataProcessor; @Autowired private ConcurrentTaskExecutor taskExecutor; @Autowired private JobStatsHolder statsHolder; // 每日调度执行 public void execute() { // 重置统计,避免上次执行的数据干扰 statsHolder.reset(); // 分批读取数据的逻辑 boolean hasMoreData = true; while (hasMoreData) { List<DataItem> batch = reader.readBatch(100); // 每次读100条 if (batch.isEmpty()) { hasMoreData = false; break; } // 提交到线程池处理 batch.forEach(item -> taskExecutor.execute(() -> dataProcessor.processData(item))); } // 等待所有线程池任务完成(关键!否则统计会提前输出) waitForTaskCompletion(); // 汇总打印统计 statsHolder.printStats(); } private void waitForTaskCompletion() { // 从ConcurrentTaskExecutor中获取内部线程池,等待任务结束 if (taskExecutor.getThreadPoolExecutor() instanceof ThreadPoolExecutor) { ThreadPoolExecutor executor = (ThreadPoolExecutor) taskExecutor.getThreadPoolExecutor(); executor.shutdown(); try { // 根据任务实际耗时设置超时时间 if (!executor.awaitTermination(1, TimeUnit.HOURS)) { executor.shutdownNow(); } } catch (InterruptedException e) { executor.shutdownNow(); Thread.currentThread().interrupt(); } } } }
方案二:用TaskDecorator增强线程池任务
如果不想侵入Reader和Processor的业务代码,可以给ConcurrentTaskExecutor添加TaskDecorator,在任务执行前后自动更新统计,这种方式更解耦。
示例代码:
@Configuration public class TaskExecutorConfig { @Bean public ConcurrentTaskExecutor concurrentTaskExecutor(JobStatsHolder statsHolder) { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(10); executor.setMaxPoolSize(20); executor.setQueueCapacity(100); // 设置任务装饰器,统一处理任务执行统计 executor.setTaskDecorator(runnable -> () -> { try { runnable.run(); statsHolder.incrementSuccessfulProcesses(); } catch (Exception e) { statsHolder.incrementFailedProcesses(); throw e; } finally { statsHolder.incrementTotalProcesses(); } }); executor.initialize(); return new ConcurrentTaskExecutor(executor); } @Bean public JobStatsHolder jobStatsHolder() { return new JobStatsHolder(); } }
这种方式下,Processor只需要专注于业务逻辑,统计逻辑由装饰器统一处理;而Reader的统计还是需要在Reader类中添加,因为读取失败是在API调用阶段,不属于线程池任务的范畴。
关键注意事项
- 线程安全:必须用
AtomicInteger、ConcurrentHashMap这类线程安全的容器,绝对不能用普通的int变量,否则并发场景下会出现统计不准确的问题。 - 等待任务完成:Job执行完读取逻辑后,必须等待线程池所有任务都处理完毕再统计,否则会漏掉还在执行的任务数据。
- 统计重置:每次Job执行前要重置统计数据,避免多次执行的统计累积。
内容的提问来源于stack exchange,提问作者JDev
相关产品推荐
相关产品推荐

