Spring Batch 5中无AsyncItemProcessor时如何实现异步Item处理与写入?
实现Spring Batch异步Item处理与写入
背景说明
AsyncItemProcessor和AsyncItemWriter在新版Spring Batch中已被移除,替代方案是利用Spring框架的异步编程特性(@Async、TaskExecutor)结合Spring Batch原生组件实现异步处理与写入。
步骤1:配置异步任务线程池
首先定义线程池用于处理异步任务,可根据业务需求调整参数:
@Configuration public class AsyncConfig { @Bean public TaskExecutor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(10); executor.setQueueCapacity(25); executor.setThreadNamePrefix("batch-async-"); executor.initialize(); return executor; } }
步骤2:改造ItemProcessor为异步实现
修改原处理器,通过@Async标记异步处理方法,返回CompletableFuture类型包装处理结果:
public class ItemProcessorImpl implements ItemProcessor<Employee, CompletableFuture<Employee>> { @Async("taskExecutor") @Override public CompletableFuture<Employee> process(Employee employee) throws Exception { // 原有业务处理逻辑,示例中模拟耗时操作 Thread.sleep(100); employee.setProcessed(true); return CompletableFuture.completedFuture(employee); } }
步骤3:添加Future结果解析处理器
新增处理器用于等待异步任务完成,提取最终处理后的对象:
public class FutureUnwrapProcessor implements ItemProcessor<CompletableFuture<Employee>, Employee> { @Override public Employee process(CompletableFuture<Employee> futureEmployee) throws Exception { return futureEmployee.get(); // 阻塞等待异步任务完成并获取结果 } }
步骤4:更新Item配置类
调整ItemConfig,注册异步处理器和结果解析处理器:
@Configuration public class ItemConfig { @Autowired private EntityManagerFactory entityManagerFactory; @Bean public JpaPagingItemReader<Employee> itemReader() { JpaPagingItemReaderBuilder<Employee> reader = new JpaPagingItemReaderBuilder<>(); reader.entityManagerFactory(entityManagerFactory); reader.pageSize(1000); reader.queryString("SELECT e FROM employees e"); reader.name("itemReader"); return reader.build(); } @Bean public ItemProcessor<Employee, CompletableFuture<Employee>> asyncProcessor() { return new ItemProcessorImpl(); } @Bean public ItemProcessor<CompletableFuture<Employee>, Employee> unwrapProcessor() { return new FutureUnwrapProcessor(); } @Bean public ItemWriter<Employee> itemWriter() { final JpaItemWriter<Employee> writer = new JpaItemWriter<>(); writer.setEntityManagerFactory(entityManagerFactory); return writer; } }
步骤5:调整Step配置
修改BatchConfig,串联异步处理器与结果解析器,完成异步处理链路:
@Configuration public class BatchConfig { @Bean public Job job(JobRepository jobRepository, Step step) { return new JobBuilder("job", jobRepository) .start(step) .build(); } @Bean public Step step(JobRepository jobRepository, PlatformTransactionManager transactionManager, JpaPagingItemReader<Employee> itemReader, ItemProcessor<Employee, CompletableFuture<Employee>> asyncProcessor, ItemProcessor<CompletableFuture<Employee>, Employee> unwrapProcessor, ItemWriter<Employee> itemWriter) { return new StepBuilder("step", jobRepository) .<Employee, Employee>chunk(1000, transactionManager) .reader(itemReader) .processor(asyncProcessor) .processor(unwrapProcessor) .writer(itemWriter) .build(); } }
可选:异步写入实现
若需要异步写入,可包装原ItemWriter实现异步提交:
public class AsyncItemWriter implements ItemWriter<Employee> { private final ItemWriter<Employee> delegate; private final TaskExecutor taskExecutor; public AsyncItemWriter(ItemWriter<Employee> delegate, TaskExecutor taskExecutor) { this.delegate = delegate; this.taskExecutor = taskExecutor; } @Override public void write(List<? extends Employee> items) throws Exception { taskExecutor.execute(() -> { try { delegate.write(items); } catch (Exception e) { throw new RuntimeException("异步写入失败", e); } }); } }
修改ItemConfig中的writer配置,替换为异步包装类:
@Bean public ItemWriter<Employee> itemWriter(EntityManagerFactory entityManagerFactory, TaskExecutor taskExecutor) { final JpaItemWriter<Employee> writer = new JpaItemWriter<>(); writer.setEntityManagerFactory(entityManagerFactory); return new AsyncItemWriter(writer, taskExecutor); }
关键注意事项
- 事务隔离:异步操作脱离原chunk事务,若需事务支持,需在异步方法上单独添加
@Transactional,此时每个异步任务将拥有独立事务。 - 异常处理:需在
Future.get()或异步任务中捕获并处理异常,避免未处理异常导致批量任务中断。 - 线程池调优:根据任务耗时、系统资源调整线程池参数,避免资源耗尽或线程闲置。
内容的提问来源于stack exchange,提问作者Renan Geraldo
相关产品推荐
相关产品推荐

