Spring Batch任务Web服务调用并行化配置及问题求助
Spring Batch任务并行化实现问题
我有一个定时调度的Spring Batch任务,从PostgreSQL数据库读取记录。针对每条记录,需要调用一次Web服务并将响应存入数据库。每次Web服务调用约需3秒响应时间,希望通过并行化提升效率。
请问如何实现该并行化?
代码片段
配置类AsyncJobConfig.java
@Configuration @EnableBatchProcessing @ConditionalOnProperty(value="fiscal.operation.mode", havingValue="async") public class AsyncJobConfig { @Autowired DocumentEntityItemWriterAsync documentEntityItemWriterAsync; @Autowired AsyncProcessor asyncProcessor; @Autowired private DocumentService documentService; @Autowired private EntityManagerFactory entityManagerFactory; @Autowired private JobBuilderFactory jobBuilderFactoryStatus; @Autowired private StepBuilderFactory stepBuilderFactoryStatus; @Autowired private DocumentEntityRepository repository; Logger logger = LogManager.getLogger(AsyncJobConfig.class); @Bean public TaskExecutor asyncTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(Integer.parseInt(System.getenv("MAX_THREAD_GET_STATUS"))); // 设置使用的线程数 executor.setMaxPoolSize(Integer.parseInt(System.getenv("MAX_THREAD_GET_STATUS"))); // 设置最大线程数 executor.setQueueCapacity(Integer.parseInt(System.getenv("MAX_THREAD_GET_STATUS"))); // 设置等待执行的任务队列容量 executor.initialize(); return executor; } @Bean public Job statusJob() { return jobBuilderFactoryStatus .get("statusJob") .incrementer(new RunIdIncrementer()) .start(initialDelayStep()) // 使用自定义步骤实现初始延迟 .next(statusStep()) // 继续执行实际处理步骤 .build(); } @Bean public Step initialDelayStep() { return stepBuilderFactoryStatus .get("initialDelayStep") .tasklet((contribution, chunkContext) -> { Thread.sleep(Long.parseLong(System.getenv("MAX_SLEEP_ASYNC"))); return RepeatStatus.FINISHED; }) .build(); } @Bean public Step statusStep() { return stepBuilderFactoryStatus .get("statusStep") .<DocumentEntity, DocumentEntity>chunk(20) .reader(statusReader()) .processor(asyncProcessor) .writer(documentEntityItemWriterAsync) .taskExecutor(asyncTaskExecutor()) // 设置任务执行器 .throttleLimit(Integer.parseInt(System.getenv("MAX_THREAD_GET_STATUS"))) // 设置并发线程数 .build(); } @Bean public ItemReader<DocumentEntity> statusReader() { logger.info("Starting statusReader, async mode"); JpaPagingItemReader<DocumentEntity> reader = new DocumentItemReaderAsync(repository); reader.setEntityManagerFactory(entityManagerFactory); reader.setQueryString("select d from DocumentEntity d WHERE d.status = :status "); reader.setParameterValues(Collections.singletonMap("status", TBE_STATUS)); return reader; } }
读取器类DocumentItemReaderAsync.java
@StepScope @Component public class DocumentItemReaderAsync extends JpaPagingItemReader< DocumentEntity> { final private List< DocumentEntity> dataList=Collections.synchronizedList(new ArrayList<>()); private String status; private int index; @Autowired public DocumentItemReaderAsync( @Autowired DocumentEntityRepository repository) { this.index = 0; List<DocumentEntity> d = null; try { d = this.read0(repository); dataList.addAll(d); } catch (Exception e) { logger.error("Distinct query Error",e); throw new RuntimeException(e); } } private List<DocumentEntity> read0 (DocumentEntityRepository repository) { return repository.findByStatusTBE(TBE_STATUS) ; } public String getStatus() { return status; } public void setStatus(String status) { this.status = status; } @Override public DocumentEntity read() throws Exception { DocumentEntity documentEntity = null; synchronized(dataList) { //dataList = repository.myFindDistinct(); // DocumentEntity documentEntity = null; if( index < this.dataList.size()) documentEntity = dataList.get(index++); while (index < dataList.size() && dataList.get(index).getDocumentId().equals(documentEntity.getDocumentId())) { this.index++; } if (documentEntity == null) { dataList.clear(); this.index = 0; } } return documentEntity; } }
遇到的问题
- 查询仅首次执行,后续任务启动时因无记录被处理而不再执行查询;
- 若每次任务启动时查询都执行,所有记录会被重复分配给所有已开启的线程;
- 使用原生
distinct查询后,第二个问题仍存在。
解决方案
问题1:查询仅首次执行的修复
你的DocumentItemReaderAsync在构造函数中一次性加载所有数据到dataList,且@Component与@StepScope误用导致reader实例被缓存,后续任务复用实例时dataList已读取完毕,不会重新查询。
修复步骤:
- 移除
DocumentItemReaderAsync上的@Component注解,因为statusReader()方法已通过@Bean创建实例,无需组件扫描; - 将查询逻辑移至
open方法,Spring Batch会在每次Step启动时调用该方法初始化:
@StepScope public class DocumentItemReaderAsync extends JpaPagingItemReader<DocumentEntity> { private final DocumentEntityRepository repository; final private List<DocumentEntity> dataList = Collections.synchronizedList(new ArrayList<>()); private int index; public DocumentItemReaderAsync(DocumentEntityRepository repository) { this.repository = repository; this.index = 0; } @Override public void open(ExecutionContext executionContext) throws ItemStreamException { super.open(executionContext); // 每次Step启动时重新查询数据 dataList.clear(); index = 0; try { List<DocumentEntity> d = repository.findByStatusTBE(TBE_STATUS); dataList.addAll(d); } catch (Exception e) { logger.error("Distinct query Error", e); throw new RuntimeException(e); } } @Override public DocumentEntity read() throws Exception { DocumentEntity documentEntity = null; synchronized(dataList) { if( index < this.dataList.size()) { documentEntity = dataList.get(index++); while (index < dataList.size() && dataList.get(index).getDocumentId().equals(documentEntity.getDocumentId())) { this.index++; } } if (documentEntity == null) { dataList.clear(); this.index = 0; } } return documentEntity; } }
同时修改AsyncJobConfig中的reader创建逻辑,添加@StepScope确保每次Step启动生成新实例:
@Bean @StepScope public ItemReader<DocumentEntity> statusReader(DocumentEntityRepository repository) { logger.info("Starting statusReader, async mode"); DocumentItemReaderAsync reader = new DocumentItemReaderAsync(repository); reader.setEntityManagerFactory(entityManagerFactory); return reader; }
问题2:记录重复分配的修复
多线程环境下,自定义reader的index变量非线程安全,即使加了synchronized仍存在竞态条件,导致多个线程读取同一条记录。
推荐修复方案:使用Spring Batch内置的JpaPagingItemReader,它原生支持分页,多线程下每个线程读取不同分页数据,避免重复:
@Bean @StepScope public ItemReader<DocumentEntity> statusReader() { logger.info("Starting statusReader, async mode"); JpaPagingItemReader<DocumentEntity> reader = new JpaPagingItemReader<>(); reader.setEntityManagerFactory(entityManagerFactory); reader.setQueryString("SELECT DISTINCT d FROM DocumentEntity d WHERE d.status = :status"); reader.setParameterValues(Collections.singletonMap("status", TBE_STATUS)); reader.setPageSize(20); // 与chunk size保持一致 return reader; }
如果必须自定义reader,将index改为原子类保证线程安全:
private final AtomicInteger index = new AtomicInteger(0); @Override public DocumentEntity read() throws Exception { DocumentEntity documentEntity = null; synchronized(dataList) { int currentIndex = index.getAndIncrement(); if( currentIndex < this.dataList.size()) { documentEntity = dataList.get(currentIndex); while (index.get() < dataList.size() && dataList.get(index.get()).getDocumentId().equals(documentEntity.getDocumentId())) { index.getAndIncrement(); } } if (documentEntity == null) { dataList.clear(); index.set(0); } } return documentEntity; }
问题3:distinct查询仍重复的修复
出现该问题的原因通常是:
DocumentEntity未正确重写equals和hashCode方法,JPA无法判断重复;- 多线程下分页逻辑未正确隔离。
修复步骤:
- 为
DocumentEntity重写equals和hashCode,基于唯一标识字段(如documentId):
@Entity public class DocumentEntity { @Id private Long documentId; // 其他字段 @Override public boolean equals(Object o) { if (this == o) return true; if (o == null || getClass() != o.getClass()) return false; DocumentEntity that = (DocumentEntity) o; return Objects.equals(documentId, that.documentId); } @Override public int hashCode() { return Objects.hash(documentId); } }
- 使用内置
JpaPagingItemReader,它会为每个线程维护独立分页状态,避免重复读取。
额外优化建议
- 线程池配置:当前线程池的
corePoolSize、maxPoolSize、queueCapacity设为同一值,建议调整queueCapacity为更大值,避免任务被拒绝; - 记录状态锁:查询时使用
SELECT ... FOR UPDATE SKIP LOCKED语句,避免多线程同时读取同一条记录:SELECT d FROM DocumentEntity d WHERE d.status = :status FOR UPDATE SKIP LOCKED - 异步处理器:如果Web服务调用是阻塞的,可给
AsyncProcessor添加@Async注解实现异步调用,同时确保处理器线程安全。
内容的提问来源于stack exchange,提问作者SoxoS
相关产品推荐
相关产品推荐

