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

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;
    }

}

遇到的问题

  1. 查询仅首次执行,后续任务启动时因无记录被处理而不再执行查询;
  2. 若每次任务启动时查询都执行,所有记录会被重复分配给所有已开启的线程;
  3. 使用原生distinct查询后,第二个问题仍存在。

解决方案

问题1:查询仅首次执行的修复

你的DocumentItemReaderAsync在构造函数中一次性加载所有数据到dataList,且@Component与@StepScope误用导致reader实例被缓存,后续任务复用实例时dataList已读取完毕,不会重新查询。

修复步骤:

  1. 移除DocumentItemReaderAsync上的@Component注解,因为statusReader()方法已通过@Bean创建实例,无需组件扫描;
  2. 将查询逻辑移至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查询仍重复的修复

出现该问题的原因通常是:

  1. DocumentEntity未正确重写equals和hashCode方法,JPA无法判断重复;
  2. 多线程下分页逻辑未正确隔离。

修复步骤:

  1. 为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);
    }
}
  1. 使用内置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 18:35:18