Spring Batch读取CSV与数据库分块处理机制咨询
Spring Batch Chunk处理机制与数据库场景实现
一、CSV文件处理的工作机制
当你设置chunk size=10处理10000行CSV时,Spring Batch是以流的方式分批次处理,具体流程如下:
- Reader每次读取1行数据,重复调用直到凑够10条(达到chunk size)
- 这10条数据被批量传给Processor,逐个完成处理
- 处理后的10条数据被批量传给Writer,一次性写入目标文件
- 重复上述步骤,直到所有10000行数据处理完成(共执行1000次chunk循环)
绝对不会一次性加载全部10000行到内存,这也是Spring Batch处理大数据量的核心优势之一。
二、Spring Batch的核心价值(为什么不用自己写简单方法)
自己写三个方法确实能完成小数据量的处理,但面对生产场景,Spring Batch的优势会凸显:
- 事务安全:每个chunk是一个独立事务,一旦处理失败,只会回滚当前chunk的操作,不会丢失之前已完成的处理结果
- 重启恢复:任务中途崩溃后,下次启动可以从失败的chunk位置继续执行,无需从头开始
- 错误处理:内置跳过(Skip)、重试(Retry)机制,可配置跳过无效数据或重试失败操作,不用手动编写复杂的异常处理逻辑
- 监控与审计:自动记录任务执行状态、处理条数、耗时等数据,方便排查问题和审计
- 扩展性:支持并行处理、分区处理,数据量极大时可以横向扩展提升处理效率
三、数据库间批量迁移的实现(分批次读/写)
完全支持每次读取10条、处理后写入数据库,以下是核心配置伪代码:
@Configuration public class DbToDbBatchJob { // 从数据库A读取数据的Reader,自动分页每次取10条 @Bean public JdbcCursorItemReader<User> sourceDbReader(DataSource sourceDataSource) { String query = "SELECT id, name, email FROM user"; return new JdbcCursorItemReaderBuilder<User>() .dataSource(sourceDataSource) .sql(query) .rowMapper((rs, rowNum) -> { User user = new User(); user.setId(rs.getLong("id")); user.setName(rs.getString("name")); user.setEmail(rs.getString("email")); return user; }) .fetchSize(10) // 数据库层面每次抓取10条数据 .build(); } // 数据处理器(按需自定义处理逻辑) @Bean public ItemProcessor<User, User> userProcessor() { return user -> { // 示例:将用户名转为大写 user.setName(user.getName().trim().toUpperCase()); return user; }; } // 写入数据库B的Writer,批量写入10条 @Bean public JdbcBatchItemWriter<User> targetDbWriter(DataSource targetDataSource) { String insertSql = "INSERT INTO user_copy(id, name, email) VALUES (:id, :name, :email)"; return new JdbcBatchItemWriterBuilder<User>() .dataSource(targetDataSource) .sql(insertSql) .beanMapped() // 自动映射User对象属性到SQL参数 .build(); } // Step配置:chunk size=10,每10条作为一个事务单元 @Bean public Step dataTransferStep(JobRepository jobRepository, PlatformTransactionManager transactionManager, JdbcCursorItemReader<User> sourceDbReader, ItemProcessor<User, User> userProcessor, JdbcBatchItemWriter<User> targetDbWriter) { return new StepBuilder("dataTransferStep", jobRepository) .<User, User>chunk(10, transactionManager) .reader(sourceDbReader) .processor(userProcessor) .writer(targetDbWriter) .build(); } // Job配置 @Bean public Job dataTransferJob(JobRepository jobRepository, Step dataTransferStep) { return new JobBuilder("dataTransferJob", jobRepository) .incrementer(new RunIdIncrementer()) .start(dataTransferStep) .build(); } }
注:需要配置两个不同的
DataSource分别对应数据库A和数据库B,Spring会自动根据Bean注入对应数据源。
内容的提问来源于stack exchange,提问作者Hua Liang
相关产品推荐
相关产品推荐

