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

如何让Spring Batch中自定义Reader的Chunk传递给自定义Writer的write方法?

Spring Batch 规范实现指导:自动分块与Writer调用问题修复

当前代码的核心问题

你现在的实现完全绕过了Spring Batch的分块机制:

  • 在reader()方法里手动打开Reader并一次性读取所有数据到全局Chunk,Spring Batch没机会按配置的chunkSize自动收集数据
  • writer()方法手动调用write,跳过了Spring Batch的事务管理和分块生命周期,数据量大时会直接内存溢出,也不符合框架的设计逻辑

规范实现步骤

  1. Reader只负责配置,不手动迭代:JdbcCursorItemReader本身就是ItemReader接口实现类,Spring Batch会自动循环调用它的read()方法,直到返回null
  2. Writer只需实现ItemWriter接口:Spring Batch会自动把分好的Chunk传递给write()方法,不需要手动处理数据收集
  3. Step配置正确利用分块机制:通过.chunk(chunkSize, transactionManager)让Spring Batch自动管理分块、事务和Reader/Writer的调用

修正后的代码

1. Reader 实现(CustomReader)

@Component
public class CustomReader {

    private final DataSource dataSource;

    @Autowired
    public CustomReader(DataSource dataSource) {
        this.dataSource = dataSource;
    }

    public JdbcCursorItemReader<MyEntity> itemReader() throws Exception {
        System.out.println("We are in the database reader");

        PreparedStatementSetter preparedStatementSetter = ps -> {
            // 填写你的参数设置逻辑
            // 示例:ps.setLong(1, targetId);
        };

        JdbcCursorItemReader<MyEntity> itemReader = new JdbcCursorItemReader<>();
        itemReader.setDataSource(dataSource);
        itemReader.setPreparedStatementSetter(preparedStatementSetter);
        itemReader.setName("myEntityReader");
        itemReader.setSql("SELECT * FROM your_source_table"); // 替换为你的查询SQL
        itemReader.setRowMapper(new CustomRowMapper());

        // 必须调用初始化方法完成配置加载
        itemReader.afterPropertiesSet();
        return itemReader;
    }
}

2. Writer 实现(CustomWriter)

@Component
public class CustomWriter implements ItemWriter<MyEntity> {

    private final NamedParameterJdbcOperations namedParameterJdbcTemplate;

    @Autowired
    public CustomWriter(DataSource dataSource) {
        this.namedParameterJdbcTemplate = new NamedParameterJdbcTemplate(dataSource);
    }

    @Override
    public void write(Chunk<? extends MyEntity> chunk) throws Exception {
        System.out.println("Hello! 正在处理Chunk,大小:" + chunk.size());

        // 批量更新逻辑:直接对整个Chunk做批量操作,无需循环每个实体
        namedParameterJdbcTemplate.getJdbcOperations().batchUpdate(
                "INSERT INTO target_table (col1, col2) VALUES (?, ?)", // 替换为你的SQL
                new BatchPreparedStatementSetter() {
                    @Override
                    public void setValues(PreparedStatement ps, int i) throws SQLException {
                        MyEntity entity = chunk.getItems().get(i);
                        // 按顺序设置SQL参数
                        ps.setString(1, entity.getCol1());
                        ps.setInt(2, entity.getCol2());
                    }

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

        // 第二个批量操作同理
        namedParameterJdbcTemplate.getJdbcOperations().batchUpdate(
                "UPDATE another_table SET status = true WHERE id = ?",
                new BatchPreparedStatementSetter() {
                    @Override
                    public void setValues(PreparedStatement ps, int i) throws SQLException {
                        MyEntity entity = chunk.getItems().get(i);
                        ps.setLong(1, entity.getId());
                    }

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

3. Batch 配置类(Config)

@Configuration
@EnableBatchProcessing
public class Config {

    private final CustomReader customReader;
    private final CustomWriter customWriter;
    private final DataSource dataSource;

    @Autowired
    public Config(DataSource dataSource, CustomReader customReader, CustomWriter customWriter) {
        this.dataSource = dataSource;
        this.customReader = customReader;
        this.customWriter = customWriter;
    }

    @Bean
    public PlatformTransactionManager transactionManager() {
        return new DataSourceTransactionManager(dataSource);
    }

    @Bean
    public ItemReader<MyEntity> myEntityReader() throws Exception {
        // 直接返回配置好的Reader,Spring Batch会自动调用其read()方法
        return customReader.itemReader();
    }

    @Bean
    public ItemWriter<MyEntity> myEntityWriter() {
        // 直接返回Writer实例,Spring Batch会自动分块调用write()
        return customWriter;
    }

    @Bean
    public Step myEntityStep(JobRepository jobRepository, PlatformTransactionManager transactionManager) throws Exception {
        return new StepBuilder("myEntityStep", jobRepository)
                .<MyEntity, MyEntity>chunk(10, transactionManager) // 指定分块大小为10,事务由框架管理
                .reader(myEntityReader())
                .writer(myEntityWriter())
                .build();
    }

    @Bean
    public Job myEntityJob(JobRepository jobRepository, Step myEntityStep) {
        return new JobBuilder("myEntityJob", jobRepository)
                .incrementer(new RunIdIncrementer())
                .start(myEntityStep)
                .build();
    }

    public static void main(String[] args) {
        System.exit(SpringApplication.exit(SpringApplication.run(Config.class, args)));
    }
}

关键说明

  • Spring Batch会自动循环调用ItemReader.read(),每收集到10条数据(chunkSize=10)就开启事务,调用ItemWriter.write()处理该Chunk,完成后提交事务
  • 若处理中出现异常,当前Chunk的事务会回滚,框架会按默认策略重试(可自定义重试/跳过规则)
  • 这种实现符合Spring Batch最佳实践,支持大数据量处理,不会一次性加载所有数据到内存,同时具备完善的事务保障

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 22:34:54