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

Spring Batch 5如何按Chunk分批读取无重叠数据并处理

问题描述

我正在使用Spring Batch 5,需要通过自定义Query按Chunk读取并处理数据。Job启动后需按Chunk(如大小为1000)分批读取数据库数据,分批处理且无数据重叠。我已编写如下代码:

1. Repository

@Transactional
@Repository
public interface TransactionRepo extends JpaRepository<TransactionEntity, Long> {
    @Query(value = "SELECT * FROM TRANSACTION WHERE STATUS =1 FOR UPDATE ", nativeQuery = true)
    List<TransactionEntity> getbKashTransaction();
}

2. Reader

@Slf4j
public class ReaderStep1 implements ItemReader<TransactionEntity> {
    private Iterator<TransactionEntity> iterator;
    private TransactionSPcall transactionSPcall;
    public ReaderStep1(TransactionRepo transactionRepo, TransactionSPcall transactionSPcall) {
        this.transactionSPcall = transactionSPcall;
        List<TransactionEntity> list = transactionRepo.getbKashTransaction();

        if(!list.isEmpty()){
            this.iterator = list.iterator();
        }
    }

    @Override
    public TransactionEntity read() {
        return (iterator != null && iterator.hasNext()) ? iterator.next() : null;
    }
}

3. Processor

@Slf4j
public class ProcessorStep1 implements ItemProcessor<TransactionEntity, TransactionEntity> {
    @Override
    public TransactionEntity process(TransactionEntity item) throws Exception {
        return item;
    }
}

4. Writer

@Slf4j
public class WriterStep1 implements ItemWriter<TransactionEntity> {
    @Override
    public void write(Chunk<? extends TransactionEntity> chunk) {
        chunk.forEach(item -> {
        });      
    }
}

5. Configuration

@Configuration
public class BatchBean {
    @Bean
    public ProcessorStep1 batchProcessorStep1() {
        return new ProcessorStep1();
    }
    @Bean
    public WriterStep1 batchWriterStep1() {
        return new WriterStep1();
    }
}

6. Job Configuration

@Component
@RequiredArgsConstructor
@Slf4j
public class BksashJob {
    private final TransactionRepo transactionRepo;
    private final ProcessorStep1 batchProcessorStep1;
    private final WriterStep1 batchWriterStep1;
    private final PlatformTransactionManager transactionManager;
    private final JobRepository jobRepository;

    @Autowired
    private TransactionSPcall transactionSPcall;

    public Job createStp1(){
        return new JobBuilder("BK-STEP1", jobRepository)
                .incrementer(new RunIdIncrementer())
                .start(step1())
                .build();
    }

    private Step step1(){
        return new StepBuilder("step1", jobRepository)
                .<TransactionEntity, TransactionEntity> chunk(1000, transactionManager)
                .reader(new ReaderStep1(transactionRepo, transactionSPcall))
                .processor(batchProcessorStep1)
                .writer(batchWriterStep1)
                .build();
    }
}

请问需要修改哪些代码,才能实现按Chunk读取无重叠数据的需求?


解决方案

你当前的实现是一次性把所有符合条件的数据加载到内存,再在内存中分块处理,既不符合真正的数据库层面Chunk分批读取,还可能引发内存溢出,同时无法保证重启或并发场景下的数据不重复读取。要实现目标,需从以下几点修改:

1. 替换自定义Reader为分页式ItemReader

Spring Batch提供了JpaPagingItemReader,专门用于数据库分页读取,自动按指定Chunk大小分批查询数据,避免一次性加载全量数据。

修改Job配置中的Step,替换原有自定义Reader:

@Component
@RequiredArgsConstructor
@Slf4j
public class BksashJob {
    // 新增注入EntityManagerFactory
    private final EntityManagerFactory entityManagerFactory;
    private final TransactionRepo transactionRepo;
    private final ProcessorStep1 batchProcessorStep1;
    private final WriterStep1 batchWriterStep1;
    private final PlatformTransactionManager transactionManager;
    private final JobRepository jobRepository;

    public Job createStp1(){
        return new JobBuilder("BK-STEP1", jobRepository)
                .incrementer(new RunIdIncrementer())
                .start(step1())
                .build();
    }

    private Step step1(){
        // 构建JpaPagingItemReader
        JpaPagingItemReader<TransactionEntity> reader = new JpaPagingItemReader<>();
        reader.setEntityManagerFactory(entityManagerFactory);
        // 自定义查询语句,筛选未处理数据
        reader.setQueryString("SELECT t FROM TransactionEntity t WHERE t.status = 1");
        reader.setPageSize(1000); // 与Chunk大小保持一致
        // 设置排序规则,保证分页读取顺序稳定,避免数据重复/遗漏
        Map<String, Sort.Direction> sortMap = new HashMap<>();
        sortMap.put("id", Sort.Direction.ASC);
        reader.setSort(sortMap);

        return new StepBuilder("step1", jobRepository)
                .<TransactionEntity, TransactionEntity>chunk(1000, transactionManager)
                .reader(reader)
                .processor(batchProcessorStep1)
                .writer(batchWriterStep1)
                .build();
    }
}

2. 处理后标记数据状态,保证无重叠

分页读取能实现分批获取,但必须在处理完成后更新数据状态,避免下次读取到已处理的数据。

修改Writer,添加状态更新逻辑:

@Slf4j
public class WriterStep1 implements ItemWriter<TransactionEntity> {
    private final TransactionRepo transactionRepo;

    // 注入Repository
    public WriterStep1(TransactionRepo transactionRepo) {
        this.transactionRepo = transactionRepo;
    }

    @Override
    public void write(Chunk<? extends TransactionEntity> chunk) {
        // 批量更新更高效,避免单条save的性能损耗
        List<Long> ids = chunk.getItems().stream()
                .map(TransactionEntity::getId)
                .toList();
        transactionRepo.updateStatusByIds(2, ids); // 将状态改为已处理(比如2)
    }
}

同时在Repository中添加批量更新方法:

@Transactional
@Repository
public interface TransactionRepo extends JpaRepository<TransactionEntity, Long> {
    @Modifying
    @Query("UPDATE TransactionEntity t SET t.status = :status WHERE t.id IN :ids")
    void updateStatusByIds(@Param("status") Integer status, @Param("ids") List<Long> ids);
}

3. 移除无用代码

原来的ReaderStep1类已被JpaPagingItemReader替代,可直接删除。

4. 可选:优化锁机制与分页方式

  • 锁机制:如果存在并发处理场景,可在查询语句中添加FOR UPDATE锁定行,避免数据冲突:
    reader.setQueryString("SELECT t FROM TransactionEntity t WHERE t.status = 1 FOR UPDATE");
    
  • 键集分页:如果数据量极大,偏移量分页性能不佳,可改用键集分页(以上一批最后一条数据的ID作为下一批起始条件),避免分页偏移问题。示例如下:
    // 自定义键集分页Reader
    public class KeysetPagingReader implements ItemReader<TransactionEntity> {
        private final TransactionRepo transactionRepo;
        private Long lastId = 0L;
        private final int pageSize = 1000;
        private Iterator<TransactionEntity> currentBatch;
    
        public KeysetPagingReader(TransactionRepo transactionRepo) {
            this.transactionRepo = transactionRepo;
            loadNextBatch();
        }
    
        private void loadNextBatch() {
            List<TransactionEntity> batch = transactionRepo.findByStatusAndIdGreaterThan(1, lastId, PageRequest.of(0, pageSize));
            if (!batch.isEmpty()) {
                lastId = batch.get(batch.size() - 1).getId();
                currentBatch = batch.iterator();
            } else {
                currentBatch = Collections.emptyIterator();
            }
        }
    
        @Override
        public TransactionEntity read() {
            if (currentBatch != null && currentBatch.hasNext()) {
                return currentBatch.next();
            } else if (!currentBatch.hasNext()) {
                loadNextBatch();
                return currentBatch.hasNext() ? currentBatch.next() : null;
            }
            return null;
        }
    }
    
    对应的Repository方法:
    List<TransactionEntity> findByStatusAndIdGreaterThan(Integer status, Long id, Pageable pageable);
    

内容的提问来源于stack exchange,提问作者Enamul Haque

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 22:12:08