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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 06:39:52