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

如何使用Spring Batch循环调用存储过程直至返回0行并处理结果?

实现Spring Batch中多次调用存储过程直至返回0行的方案

核心思路

通过自定义ItemReader循环调用存储过程,每次获取一批数据;当存储过程返回空结果集时,Reader返回null触发Spring Batch停止读取。每次读取的批次数据会自动流转到ItemProcessor和ItemWriter处理,存储过程自行维护偏移量,无需Reader介入分页逻辑。

步骤实现

1. 自定义存储过程调用的ItemReader

实现ItemReader<List<YourEntity>>,每次调用存储过程并将结果映射为实体列表,空列表时返回null标识读取结束:

@Component
public class StoredProcedureBatchReader implements ItemReader<List<YourEntity>> {

    private final JdbcTemplate jdbcTemplate;

    public StoredProcedureBatchReader(JdbcTemplate jdbcTemplate) {
        this.jdbcTemplate = jdbcTemplate;
    }

    @Override
    public List<YourEntity> read() throws Exception {
        // 替换为你的存储过程调用语句
        List<YourEntity> batchData = jdbcTemplate.query(
            "CALL get_processing_batch()",
            (rs, rowNum) -> {
                YourEntity entity = new YourEntity();
                entity.setId(rs.getLong("id"));
                entity.setContent(rs.getString("content"));
                // 映射其他业务字段
                return entity;
            }
        );

        // 空结果集时返回null,告知Spring Batch读取结束
        return batchData.isEmpty() ? null : batchData;
    }
}

2. 实现批次处理的ItemProcessor

由于Reader返回的是批次列表,Processor需要接收List<YourEntity>并转换为处理后的实体列表:

@Component
public class BatchItemProcessor implements ItemProcessor<List<YourEntity>, List<ProcessedEntity>> {

    @Override
    public List<ProcessedEntity> process(List<YourEntity> rawBatch) throws Exception {
        return rawBatch.stream()
            .map(raw -> {
                ProcessedEntity processed = new ProcessedEntity();
                processed.setId(raw.getId());
                processed.setProcessedContent("HANDLED_" + raw.getContent());
                // 添加自定义业务处理逻辑
                return processed;
            })
            .collect(Collectors.toList());
    }
}

3. 实现批次写入的ItemWriter

Writer接收Processor输出的批次列表,注意参数为List<? extends List<ProcessedEntity>>(因为Chunk每次处理一个Reader返回的批次列表,会传递多个这样的列表),可扁平化后批量写入:

@Component
public class BatchItemWriter implements ItemWriter<List<ProcessedEntity>> {

    private final JdbcTemplate jdbcTemplate;

    public BatchItemWriter(JdbcTemplate jdbcTemplate) {
        this.jdbcTemplate = jdbcTemplate;
    }

    @Override
    public void write(List<? extends List<ProcessedEntity>> batchLists) throws Exception {
        // 扁平化所有批次数据
        List<ProcessedEntity> allProcessedItems = batchLists.stream()
            .flatMap(List::stream)
            .collect(Collectors.toList());

        // 批量写入目标表
        String insertSql = "INSERT INTO processed_data (id, processed_content) VALUES (?, ?)";
        jdbcTemplate.batchUpdate(insertSql, new BatchPreparedStatementSetter() {
            @Override
            public void setValues(PreparedStatement ps, int index) throws SQLException {
                ProcessedEntity item = allProcessedItems.get(index);
                ps.setLong(1, item.getId());
                ps.setString(2, item.getProcessedContent());
            }

            @Override
            public int getBatchSize() {
                return allProcessedItems.size();
            }
        });
    }
}

4. 配置Job和Step

Step中设置Chunk大小为1(因为Reader每次返回一个完整批次,Chunk处理的是这个批次对象),关联Reader、Processor和Writer:

@Configuration
@EnableBatchProcessing
public class BatchJobConfig {

    private final JobBuilderFactory jobBuilderFactory;
    private final StepBuilderFactory stepBuilderFactory;
    private final StoredProcedureBatchReader batchReader;
    private final BatchItemProcessor batchProcessor;
    private final BatchItemWriter batchWriter;

    public BatchJobConfig(JobBuilderFactory jobBuilderFactory, StepBuilderFactory stepBuilderFactory,
                          StoredProcedureBatchReader batchReader, BatchItemProcessor batchProcessor,
                          BatchItemWriter batchWriter) {
        this.jobBuilderFactory = jobBuilderFactory;
        this.stepBuilderFactory = stepBuilderFactory;
        this.batchReader = batchReader;
        this.batchProcessor = batchProcessor;
        this.batchWriter = batchWriter;
    }

    @Bean
    public Step processingStep() {
        return stepBuilderFactory.get("processingStep")
            .<List<YourEntity>, List<ProcessedEntity>>chunk(1)
            .reader(batchReader)
            .processor(batchProcessor)
            .writer(batchWriter)
            .build();
    }

    @Bean
    public Job dataProcessingJob() {
        return jobBuilderFactory.get("dataProcessingJob")
            .start(processingStep())
            .build();
    }
}

关键说明

  • Reader终止逻辑:当存储过程返回0行时,Reader返回null,Spring Batch会自动终止读取流程。
  • Chunk大小设置:因为Reader每次返回的是一个完整批次,所以Chunk size设为1,表示每次处理一个批次对象。
  • 存储过程偏移量:存储过程自行维护偏移量(比如标记已处理数据或维护偏移量表),Reader无需关心分页逻辑,只需重复调用即可。
  • 异常处理:可在Reader中添加异常捕获逻辑,结合Spring Batch的重试机制处理存储过程调用失败的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 14:03:26