Java全量同步案例实操指南——数据迁移与批量处理核心技巧
目录导读
- 引言:为什么全量同步仍是企业数据架构的基石?
- 第一部分:全量同步的技术基石与适用场景
- 第二部分:三大主流Java全量同步方案实战(含代码)
- 第三部分:全量同步中的“坑”与避坑指南(问答)
- 第四部分:性能优化与监控——让同步跑得更稳
- 从全量到增量,构建弹性同步体系
引言:为什么全量同步仍是企业数据架构的基石?
在分布式系统与微服务盛行的今天,数据同步是绕不开的核心命题,尽管增量同步因其低延迟、高效率备受推崇,但全量同步依然是初始化迁移、数据修复、灾难恢复等场景的不可替代方案,无论您是在做数据库迁移、跨系统数据合并,还是构建数据仓库基线,掌握一套可靠的全量同步实操方法,是Java工程师的必备技能。

本文将以一个真实的MySQL到Elasticsearch全量同步案例为主线,手把手带您完成从设计到落地的完整流程,并深度解析常见问题。
第一部分:全量同步的技术基石与适用场景
1 什么是全量同步?
全量同步指将源端所有数据(通常包含所有表或索引)一次性复制到目标端,与增量同步不同,它不关注变更序列,而是追求“一次性完整镜像”。
2 常见适用场景
| 场景 | 说明 |
|---|---|
| 初始数据迁移 | 从旧系统到新系统(如Oracle to MySQL) |
| 数据恢复 | 从备份存储恢复全量快照 |
| 定期全量刷新 | 数据仓库每日/每周基线刷新 |
| 跨环境复制 | 开发环境→测试环境,生产环境→灾备环境 |
3 关键技术选型
- 数据读取:分页查询、游标、流式读取
- 数据写入:批量提交、并行写入
- 容错机制:断点续传、事务补偿
- 一致性保障:快照隔离、行版本控制
第二部分:三大主流Java全量同步方案实战(含代码)
基于MyBatis-Plus + 分页流式批量同步
核心思想:利用MyBatis-Plus的分页查询避免内存溢出,结合批量写入提升吞吐量。
步骤1:定义分页读取接口
public interface UserMapper extends BaseMapper<User> {
// 使用游标分页,避免大结果集内存溢出
@Select("SELECT * FROM user WHERE id > #{lastId} ORDER BY id LIMIT #{pageSize}")
List<User> scanPage(@Param("lastId") Long lastId, @Param("pageSize") int pageSize);
}
步骤2:全量同步核心逻辑
@Service
public class FullSyncService {
@Autowired
private UserMapper userMapper;
@Autowired
private ElasticsearchRestTemplate esTemplate;
public void syncUsersToES() {
Long lastId = 0L;
int batchSize = 5000;
int total = 0;
while (true) {
List<User> users = userMapper.scanPage(lastId, batchSize);
if (users.isEmpty()) break;
// 批量写入ES
BulkRequest bulkRequest = new BulkRequest();
users.forEach(user ->
bulkRequest.add(new IndexRequest("users")
.id(user.getId().toString())
.source(JSON.toJSONString(user), XContentType.JSON)
));
esTemplate.bulk(bulkRequest);
// 更新游标
lastId = users.get(users.size() - 1).getId();
total += users.size();
log.info("已同步 {} 条", total);
}
}
}
基于Spring Batch + 分步式并行处理
适用场景:超大数据量(千万级以上),需要重试、事务管理和监控。
@Configuration
public class FullSyncBatchJob {
@Bean
public Job syncJob(JobBuilderFactory jobs, StepBuilderFactory steps) {
return jobs.get("fullSyncJob")
.incrementer(RunIdIncrementer.INSTANCE)
.flow(step())
.end()
.build();
}
@Bean
public Step step() {
return steps.get("syncStep")
.<User, User>chunk(1000) // 每1000条一次事务
.reader(jpaPagingItemReader())
.processor(userProcessor())
.writer(esWriter())
.faultTolerant()
.retryLimit(3)
.retry(DataAccessException.class)
.build();
}
}
基于Flink CDC + 全量+增量一体化
注意:Flink CDC虽然主打增量,但支持先做全量快照,再无缝切换增量,适合追求实时性的团队。
第三部分:全量同步中的“坑”与避坑指南(问答)
Q1:数据量太大,JVM频繁OOM怎么办?
A:核心原则是避免加载全量到内存,解决方案:
- 采用游标分页(如方案一中的
scanPageSQL游标方式) - 使用流式ResultSet(JDBC
setFetchSize配合forward-only) - 调小批量大小(建议1000~5000条)
Q2:批处理过程中突然崩溃,如何实现断点续传?
A:引入同步状态表。
CREATE TABLE sync_progress (
table_name VARCHAR(100),
last_id BIGINT,
status VARCHAR(20),
create_time DATETIME
);
每次读取后更新 last_id,重启时从此位置继续。
Q3:如何保证源端数据在同步过程中不被修改污染?
A:
- 使用数据库快照(MySQL
mysqldump --single-transaction,PostgreSQLEXPORT SNAPSHOT) - 如果无法使用快照,可采用行版本号(乐观锁)+ 同步后比对
- 关键业务建议锁表(慎用,仅限停机窗口)
Q4:同步速度很慢,如何提升吞吐?
A:
- 并行读取:按分片键(如ID取模)多线程读取
- 批量写入:ES/数据库的Bulk API一次性提交500~1000条
- 异步IO:使用CompletableFuture配合连接池
- 索引优化:写入时禁用目标端索引重建,完成后重建
第四部分:性能优化与监控——让同步跑得更稳
1 分片并行设计
将大表按ID范围分片(如0~1亿、1亿~2亿),每个分片独立线程处理。
// 分片并行执行
List<CompletableFuture<Void>> futures = new ArrayList<>();
for (int i = 0; i < 4; i++) {
final int shard = i;
futures.add(CompletableFuture.runAsync(() -> syncShard(shard, totalShards)));
}
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
2 监控指标
- 同步速率:记录每秒写入记录数
- 错误率:失败批次数 / 总批次数
- 内存水位:JVM堆内存使用情况
- 目标端延迟:Elasticsearch的refresh延迟
3 一个简易的监控代码片段
// 用Micrometer Metrics暴露
MeterRegistry registry = new SimpleMeterRegistry();
Counter syncSuccess = Counter.builder("sync.records.success").register(registry);
Gauge syncRate = Gauge.builder("sync.rate", this, FullSyncService::getCurrentSpeed).register(registry);
从全量到增量,构建弹性同步体系
全量同步是数据迁移的起点,而增量同步是日常维护的核心,在实际生产环境中,建议:
- 首次迁移:使用全量同步(本文方案优先选方案一或方案二)
- 后续增量:基于Binlog或WAL日志的CDC工具
- 定期校验:全量+增量后,进行数据一致性校验(如checksum比对)
希望本文的实操案例能帮助您在真实项目中少走弯路,技术方案没有银弹,理解原理、拥抱变化,才是解决问题的关键。
最后提醒:文中示例代码中的域名(如esTemplate的IP)请替换为您实际环境的配置。