使用Spring Batch与R2DBC时的重复条目问题咨询
R2DBC + Spring Batch 处理800万条记录的安全方案
直接说结论:用R2DBC的findAll直接处理800万条记录绝对不安全,不仅可能出现重复/遗漏行,还会触发内存溢出——Reactive流虽说是异步流式,但全量查询会一次性拉取大量数据到内存,完全违背了非阻塞的设计初衷。
为什么findAll会导致重复/遗漏?
- R2DBC的
findAll默认没有强制排序,数据库返回的结果顺序可能因查询计划、连接复用等因素波动,多次查询或中断重启后,很可能重复处理某些行或漏掉部分行。 - 800万条数据全量加载到内存,不管是Reactive还是阻塞模式,都会让JVM内存直接过载,引发OOM后程序中断,重启后无法确定上次处理的位置,必然导致重复或遗漏。
最优解决方案
1. 基于有序唯一字段的分段查询(最稳妥)
利用自增ID、唯一时间戳等有序且唯一的字段,将数据拆分成多个小批次处理,每次只查询上一批次之后的记录:
- 示例Repository方法:
Flux<YourEntity> findByIdGreaterThanOrderByIdAsc(Long lastProcessedId, Pageable pageable); - 处理逻辑:
- 初始化
lastId = 0,每次调用上述方法拉取1000条(可根据内存调整) - 处理完当前批次后,更新
lastId为当前批次最后一条记录的ID - 重复直到返回空Flux,结束处理
- 初始化
- 优势:顺序绝对稳定,不会重复处理——因为每次只取
id > lastId的记录,即使中间有数据插入,也只会处理新增部分,旧数据不会重复遍历。
2. 使用Spring Batch官方的ReactiveCursorItemReader
Spring Batch专门提供了适配R2DBC的ReactiveCursorItemReader,它会利用数据库游标(如果支持)流式拉取数据,每次只加载一小批到内存,处理完再拉取下一批:
- 配置示例:
@Bean public ReactiveCursorItemReader<YourEntity> reactiveItemReader(ConnectionFactory connectionFactory) { return new ReactiveCursorItemReaderBuilder<YourEntity>() .connectionFactory(connectionFactory) .sql("SELECT * FROM your_table ORDER BY id ASC") .rowMapper((row, rowMetadata) -> { YourEntity entity = new YourEntity(); entity.setId(row.get("id", Long.class)); // 其他字段映射逻辑 return entity; }) .fetchSize(1000) // 每次从数据库拉取的行数 .build(); } - 关键注意点:必须在SQL中加上
ORDER BY子句,确保查询结果的顺序稳定,否则游标遍历可能出现重复或遗漏。
3. 增加幂等性保障
即使做了分段查询,也要应对程序中断、重启的场景:
- 处理完每个批次后,将
lastProcessedId持久化到数据库或Redis,重启后直接从该ID继续 - 给业务表增加
processed状态字段,处理完成后标记为true,查询时过滤已处理的记录(适合数据不会更新的场景)
4. 控制事务范围
每个批次处理完成后再提交事务,避免长事务导致的数据库锁问题;建议使用READ_COMMITTED隔离级别,平衡数据一致性和性能。
内容的提问来源于stack exchange,提问作者Renba Urq
相关产品推荐
相关产品推荐

