Spring Batch实现变更检测:避免重复写入MongoDB的方案咨询
Spring Batch 700万数据变更检测最佳实践
针对你要实现的「处理后数据与上一次结果一致则不写入Mongo」的需求,先分析你提出的两个方案,再给出更贴合批量处理场景的优化建议:
方案1:自定义Reader聚合SQL与Mongo数据
优势
- 职责清晰:Reader负责数据的获取与聚合,Processor只专注于业务处理和变更判断,符合单一职责原则;
- 性能可控:可以通过批量查询Mongo替代单条查询——比如每次读取一个chunk的SQL数据后,用ID列表做
$in查询批量拉取对应历史数据,比Processor里单条查库效率提升几个数量级; - 适配多线程:如果用Spring Batch的多线程TaskExecutor,每个Reader线程独立处理自己的chunk数据,不会有线程安全问题。
注意事项
- 不要一次性加载全量数据:保持原有JDBC Reader的分页逻辑(比如
JdbcPagingItemReader),按chunk批量读取SQL数据,再批量查询Mongo,避免OOM; - 数据配对要高效:用ID作为key将Mongo查询结果存入HashMap,再和SQL数据一一映射成DTO(比如
Pair<SqlData, MongoHistoryData>)传入Processor。
方案2:Listener+内存缓存
优势
- 侵入性低:无需修改现有Reader,只需要添加监听器即可实现;
- 逻辑灵活:可以在Chunk级别做缓存,避免单条数据的重复查询。
潜在问题与优化
- 多线程安全:如果用多线程处理,必须用
ThreadLocal存储缓存(每个线程只存自己当前chunk的历史数据),或者用线程安全的ConcurrentHashMap并按线程ID隔离数据; - 内存占用:必须在
afterChunk阶段清空缓存,避免700万数据累积导致OOM; - 分布式场景限制:如果用远程分区的多节点批处理,内存缓存无法跨节点共享,这个方案就不适用了。
推荐实现方式
优先选择优化后的方案1,具体步骤如下:
- 基于现有
JdbcPagingItemReader扩展自定义Reader,保留分页读取SQL数据的逻辑; - 每次读取一个chunk的SQL数据后,提取所有数据的唯一ID,调用MongoRepository的批量查询方法(比如
findAllById(Collection<ID> ids)); - 将Mongo查询结果以ID为key存入HashMap,遍历SQL数据,将每条SQL数据与对应的Mongo历史数据包装成DTO;
- Processor接收DTO,对比新旧数据的业务字段,若未变更则返回
null(Spring Batch会自动跳过Writer写入),若变更则返回处理后的新数据; - Writer只处理Processor返回的非null数据,写入Mongo。
额外优化思路
如果业务数据有可靠的变更标识(比如SQL表中的update_time字段),可以在SQL查询阶段提前过滤:
- 记录上一次作业的结束时间,在SQL查询时只拉取
update_time > 上次作业时间的数据,从源头减少需要处理的数据量,进一步提升性能。
内容的提问来源于stack exchange,提问作者Abhi
相关产品推荐
相关产品推荐

