Spring Batch中如何在AsyncItemProcessor异步处理时更新计数器?
解决AsyncItemProcessor异步执行时的计数器线程安全问题
你的问题核心在于多线程环境下直接修改StepExecution的ExecutionContext会导致线程安全问题——因为AsyncItemProcessor用线程池并行处理每个item,多个线程同时读写ExecutionContext里的计数器,会出现计数不准确、甚至并发修改异常的情况。下面是具体的解决方案和修改后的代码:
问题分析
原来的代码中,你在process方法里直接操作stepExecution.getExecutionContext(),而StepExecution和它的ExecutionContext都不是线程安全的。当多个异步线程同时执行putInt操作时,会出现线程竞争,导致计数结果错误。
解决方案
我们需要用线程安全的原子类来维护计数器,避免直接在多线程中操作ExecutionContext。具体步骤:
- 在处理器中使用
AtomicInteger来记录处理的产品数和客户数,原子类的操作是线程安全的; - 通过
StepExecutionListener在步骤结束时,将原子计数器的值统一更新到ExecutionContext(如果需要持久化这些统计值); - 避免在异步的
process方法中直接操作StepExecution。
修改后的代码示例
1. 重构BulkExportItemProcessor,使用线程安全计数器
public class BulkExportItemProcessor implements ItemProcessor<String, String>, StepExecutionListener { @Inject public IProductRepository repository; // 线程安全的原子计数器 private final AtomicInteger processedRecords = new AtomicInteger(0); private final AtomicInteger processedCustomers = new AtomicInteger(0); @Override public void beforeStep(StepExecution stepExecution) { // 初始化计数器(支持从ExecutionContext恢复上下文) processedRecords.set(stepExecution.getExecutionContext().getInt("PROCESSED_RECORDS", 0)); processedCustomers.set(stepExecution.getExecutionContext().getInt("PROCESSED_CUST", 0)); } @Override public String process(String customerIds) { String[] customerIdsList = customerIds.split(","); int customerCount = customerIdsList.length; int productCount = 0; for (String customerId : customerIdsList) { List<IProductDTO> products = repository.getProducts(customerId, null); for (IProductDTO product : products) { productCount++; System.out.println("handling product: " + product.getProductId() + " of customer: " + customerId); } } // 原子更新计数器,避免线程竞争 processedRecords.addAndGet(productCount); processedCustomers.addAndGet(customerCount); return customerIds; } @Override public ExitStatus afterStep(StepExecution stepExecution) { // Step结束后在主线程更新ExecutionContext,线程安全 stepExecution.getExecutionContext().putInt("PROCESSED_RECORDS", processedRecords.get()); stepExecution.getExecutionContext().putInt("PROCESSED_CUST", processedCustomers.get()); return ExitStatus.COMPLETED; } }
2. 调整配置类,注册StepExecutionListener
在你的ExportBatchConfiguration中,需要把处理器注册为Step的监听器,并完善线程池初始化:
@Configuration @EnableBatchProcessing public class ExportBatchConfiguration { @Inject public JobBuilderFactory jobBuilderFactory; @Inject public StepBuilderFactory stepBuilderFactory; @Bean public FlatFileItemReader<String> reader() { return new FlatFileItemReaderBuilder<String>() .name("voltItemReader") .resource(new ClassPathResource("Customers.csv")) .build(); } @Bean public BulkExportItemProcessor delegateProcessor() { return new BulkExportItemProcessor(); } @Bean public ItemProcessor<String, String> asyncItemProcessor() { AsyncItemProcessor<String, String> asyncItemProcessor = new AsyncItemProcessor<>(); asyncItemProcessor.setDelegate(delegateProcessor()); asyncItemProcessor.setTaskExecutor(getAsyncExecutor()); return asyncItemProcessor; } @Bean(name = "asyncExecutor") public TaskExecutor getAsyncExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(100); executor.setMaxPoolSize(100); executor.setQueueCapacity(100); executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.setThreadNamePrefix("AsyncExecutor-"); executor.initialize(); // 必须初始化线程池,否则会报错 return executor; } @Bean public Step exportStep() { return stepBuilderFactory.get("ProductExportJob") .<String, String>chunk(3) .reader(reader()) .processor(asyncItemProcessor()) .listener(delegateProcessor()) // 注册Step监听器,用于统一更新计数器 .build(); } @Bean public Job job() { return jobBuilderFactory.get("ProductExportJob") .start(exportStep()) .build(); } }
关键说明
- AtomicInteger的作用:它的
addAndGet方法是原子操作,多线程环境下不会出现竞争条件,保证计数完全准确; - StepExecutionListener的时机:
afterStep方法在Step的主线程中执行,此时更新ExecutionContext是线程安全的,彻底避免了多线程直接写的问题; - 线程池初始化:ThreadPoolTaskExecutor需要显式调用
initialize(),否则会出现线程池未就绪的异常。
内容的提问来源于stack exchange,提问作者Dor
相关产品推荐
相关产品推荐

