如何让Spring Batch中自定义Reader的Chunk传递给自定义Writer的write方法?
Spring Batch 规范实现指导:自动分块与Writer调用问题修复
当前代码的核心问题
你现在的实现完全绕过了Spring Batch的分块机制:
- 在
reader()方法里手动打开Reader并一次性读取所有数据到全局Chunk,Spring Batch没机会按配置的chunkSize自动收集数据 writer()方法手动调用write,跳过了Spring Batch的事务管理和分块生命周期,数据量大时会直接内存溢出,也不符合框架的设计逻辑
规范实现步骤
- Reader只负责配置,不手动迭代:
JdbcCursorItemReader本身就是ItemReader接口实现类,Spring Batch会自动循环调用它的read()方法,直到返回null - Writer只需实现
ItemWriter接口:Spring Batch会自动把分好的Chunk传递给write()方法,不需要手动处理数据收集 - 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
相关产品推荐
相关产品推荐

