Spring Batch AsyncItemProcessor未并行处理配置问题求助
问题现象
配置AsyncItemProcessor期望并行处理含1000条记录的文件,但实际每次仅处理100条(与Chunk大小一致),且每个Chunk处理完成后出现2分钟延迟,期间仅能看到GC日志。
相关代码
配置类代码
@Configuration @EnableBatchProcessing @EnableAsync public class JobConfig { @Autowired private JobBuilderFactory jobBuilder; @Autowired private StepBuilderFactory stepBuilder; @Autowired @Qualifier("writer") private ItemWriter writer; @Bean @JobScope public ItemProcessor itemProcessor() { ItemProcessor itemProcessor = new ItemProcessor(); return itemProcessor; } @Bean @JobScope public AsyncItemProcessor asyncItemProcessor() throws IOException { AsyncItemProcessor asyncItemProcessor = new AsyncItemProcessor(); asyncItemProcessor.setDelegate(itemProcessor()); asyncItemProcessor.setTaskExecutor(getAsyncExecutor()); asyncItemProcessor.afterPropertiesSet(); return asyncItemProcessor; } @Bean(name = "asyncExecutor") public TaskExecutor getAsyncExecutor() { SimpleAsyncTaskExecutor simpleAsyncTaskExecutor = new SimpleAsyncTaskExecutor() { @Override protected void doExecute(Runnable task) { final JobExecution jobExecution = JobSynchronizationManager.getContext().getJobExecution(); super.doExecute(() -> { JobSynchronizationManager.register(jobExecution); try { task.run(); } finally { JobSynchronizationManager.close(); } }); } }; simpleAsyncTaskExecutor.setThreadNamePrefix("processing 1-"); simpleAsyncTaskExecutor.setConcurrencyLimit(100); return simpleAsyncTaskExecutor; } @Bean @JobScope public AsyncItemWriter asyncItemWriter() throws IOException { AsyncItemWriter asyncItemWriter = new AsyncItemWriter<>(); asyncItemWriter.setDelegate(writer); asyncItemWriter.afterPropertiesSet(); return asyncItemWriter; } @Bean @JobScope public FlatFileItemReader<Result> requestFileReader() { DefaultLineMapper lineMapper = new DefaultLineMapper(); ...... FlatFileItemReader<Result> itemReader = new FlatFileItemReader<>(); itemReader.setLineMapper(lineMapper); return itemReader; } @Bean public Step simpleFileStep() throws IOException { return stepBuilder.get("simpleFileStep").chunk(100).reader(fileReader).processor(asyncItemProcessor()) .writer(asyncItemWriter()).build(); } @Bean(name = "customExecutionContext") @JobScope public ExecutionContext customExecutionContext() { ExecutionContext executionContext = new ExecutionContext(); return executionContext; } }
处理器类代码
@JobScope public class RequestProcessor implements ItemProcessor<Result, List<Item>> { @Value("#{jobExecution}") private JobExecution jobExecution; @Autowired @Qualifier("customExecutionContext") private ExecutionContext storedContext; @Override public List<Item> process(Result result) throws Exception { Date start = new Date(); // Processing logic Date end = new Date(); long diff = end.getTime() - start.getTime(); log.info("Time taken to process the items:"+TimeUnit.SECONDS.convert(diff, TimeUnit.MILLISECONDS)); return items; } }
延迟期间GC日志
[GC concurrent-string-deduplication, 16.2K->0.0B(16.2K), avg 88.7%, 0.0000820 secs] [GC pause (G1 Evacuation Pause) (young) 687M->302M(768M), 0.0138859 secs]
问题分析与解决方案
1. 单次仅处理100条的原因
Spring Batch的AsyncItemProcessor是在单个Chunk内部实现并行处理,而非跨Chunk并行。你设置的chunk(100)表示每次读取100条记录,这100条会被异步处理器并行处理,处理完成后一次性写入,之后才会读取下一个Chunk。这是Chunk模型的正常行为,若要实现跨Chunk的全局并行,需要改用远程分区或分片处理方案。
2. Chunk后2分钟延迟的根源
从GC日志看,单次Young区回收耗时极短,2分钟延迟并非单次GC导致,而是大量对象堆积引发的频繁GC甚至Full GC:
- 处理器
process方法返回List<Item>,若每个Result对应返回的List数据量极大,100个此类对象会瞬间占据大量内存,Chunk处理完成后对象回收阶段会引发多次GC。 - 你设置的
concurrencyLimit(100)与Chunk大小一致,100个线程同时处理会瞬间填满年轻代,加剧内存压力。
3. 配置中的潜在问题与修复
(1) 不合理的@JobScope使用
asyncItemProcessor、asyncItemWriter、itemProcessor无需标记@JobScope,重复创建这些组件会增加不必要的开销,甚至引发线程上下文问题。移除该注解,改为单例Bean:
@Bean public ItemProcessor itemProcessor() { return new RequestProcessor(); } @Bean public AsyncItemProcessor asyncItemProcessor() throws IOException { AsyncItemProcessor asyncItemProcessor = new AsyncItemProcessor(); asyncItemProcessor.setDelegate(itemProcessor()); asyncItemProcessor.setTaskExecutor(getAsyncExecutor()); asyncItemProcessor.afterPropertiesSet(); return asyncItemProcessor; }
(2) 低效的SimpleAsyncTaskExecutor
SimpleAsyncTaskExecutor每次执行任务都会创建新线程,线程创建销毁开销极大。改用ThreadPoolTaskExecutor复用线程:
@Bean(name = "asyncExecutor") public TaskExecutor getAsyncExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(50); executor.setMaxPoolSize(100); executor.setQueueCapacity(200); executor.setThreadNamePrefix("processing-"); executor.setWaitForTasksToCompleteOnShutdown(true); executor.initialize(); // 保持JobExecution上下文传递逻辑 return task -> { JobExecution jobExecution = JobSynchronizationManager.getContext().getJobExecution(); executor.execute(() -> { JobSynchronizationManager.register(jobExecution); try { task.run(); } finally { JobSynchronizationManager.close(); } }); }; }
(3) ExecutionContext线程安全风险
异步线程共享@JobScope的ExecutionContext会引发线程安全问题,若需存储任务状态,改用线程安全容器,或每个线程独立维护状态后再合并到JobExecution上下文。
(4) 内存优化
检查处理器逻辑,减少大对象创建,及时释放无用对象;若返回的List<Item>过大,调整Writer逻辑实现分批次写入,降低单Chunk内存占用。
内容的提问来源于stack exchange,提问作者Dev

