Java案例如何实现数据迁移?从理论到实战的完整指南
目录导读
- 数据迁移的核心挑战与Java的解决之道
- 案例背景:从MySQL到PostgreSQL的跨库迁移
- 关键技术选型:Batch处理与连接池优化
- 实战代码拆解:逐行解释迁移逻辑
- 性能调优与异常处理策略
- 常见问题问答(FAQ)
数据迁移的核心挑战与Java的解决之道
在如今微服务与多数据库并存的架构中,数据迁移是开发者和运维必遇的课题,常见痛点包括:数据量大导致超时、源目标库数据类型不兼容、迁移过程中业务不断服,Java作为企业级开发的主力语言,凭借其成熟的JVM生态和丰富的数据库连接库(JDBC、MyBatis、Spring Batch),天然适合构建稳健的数据迁移工具。

迁移的本质是 “读取-转换-写入”(ETL)流水线,而Java的多线程能力与ORM框架可大幅降低编码复杂度,借助java.util.concurrent的线程池,能实现分页并发读取;通过JPA或MyBatis的流式查询,避免一次性加载全量数据导致OOM。
核心要点:数据迁移不是简单的
INSERT INTO … SELECT,必须处理脏数据、字段映射、主键冲突等问题。
案例背景:从MySQL到PostgreSQL的跨库迁移
假设某电商系统需要将订单表orders从MySQL迁移至PostgreSQL,MySQL表结构如下:
-- MySQL端 CREATE TABLE orders ( id BIGINT AUTO_INCREMENT, order_no VARCHAR(32), create_time DATETIME, amount DECIMAL(10,2), status TINYINT, PRIMARY KEY (id) );
目标PostgreSQL表结构类似,但status字段改为VARCHAR(20),且create_time需转成TIMESTAMP,核心需求:迁移历史4亿行数据,要求不停服,增量同步。
在方案设计上,我们选择分段迁移:先用Java全量迁移历史数据,再通过监听binlog实现增量同步(本案例重点讲解全量部分)。
关键技术选型:Batch处理与连接池优化
| 组件 | 选型理由 |
|---|---|
| Spring Boot 2.7 | 快速集成,自动配置数据源 |
| HikariCP | 高性能连接池,支持连接泄漏检测 |
| MyBatis-Plus | 流式查询(Cursor),避免内存溢出 |
| Java 8 Stream + Parallel | 利用多核CPU并发处理数千条/批 |
关键优化点:
- 批处理大小:MySQL JDBC写入时,设置
rewriteBatchedStatements=true,每次batch size建议1000-5000条,过大反而因事务日志激增导致性能下降。 - 连接隔离:源库与目标库使用独立HikariCP连接池,防止相互干扰。
实战代码拆解:逐行解释迁移逻辑
数据读取:游标式分页
// 使用MyBatis的Cursor避免全量加载到内存
try (Cursor<OrdersEntity> cursor = ordersMapper.scanAll()) {
List<OrdersEntity> batch = new ArrayList<>(BATCH_SIZE);
for (OrdersEntity entity : cursor) {
// 2. 数据转换:状态码转字符串,时间格式调整
entity.setStatus(convertStatus(entity.getStatus()));
entity.setCreateTime(convertTime(entity.getCreateTime()));
batch.add(entity);
if (batch.size() >= BATCH_SIZE) {
// 3. 写入目标库
postgresOrdersMapper.batchInsert(batch);
batch.clear();
}
}
// 处理剩余不足一批的数据
if (!batch.isEmpty()) postgresOrdersMapper.batchInsert(batch);
}
注意:流式查询需在application.yml禁用自动关闭连接:
mybatis.configuration.default-fetch-size: 1000 spring.datasource.hikari.read-only: true
批量写入:MyBatis的foreach优化
<!-- PostgreSQL批量插入 -->
<insert id="batchInsert" parameterType="list">
INSERT INTO orders (id, order_no, create_time, amount, status)
VALUES
<foreach collection="list" item="item" separator=",">
(#{item.id}, #{item.orderNo}, #{item.createTime}, #{item.amount}, #{item.status})
</foreach>
</insert>
避坑指南:MyBatis默认批处理不支持SELECT LAST_INSERT_ID(),若需返回自增ID,应使用useGeneratedKeys,但本案例中PostgreSQL使用序列自增,因此无需特殊处理。
异常处理:事务边界与重试机制
@Transactional(rollbackFor = Exception.class, timeout = 120)
public void batchInsert(List<OrdersEntity> list) {
try {
mapper.batchInsert(list);
} catch (DataIntegrityViolationException e) {
// 处理主键冲突:跳过该批次并记录日志
log.warn("批次存在重复数据,逐条插入: {}", e.getMessage());
list.forEach(item -> {
try { mapper.insertWithUpsert(item); }
catch (DuplicateKeyException ex) { log.error("跳过已存在记录: {}", item.getId()); }
});
}
}
性能调优与异常处理策略
调优三板斧
- 并行度控制:使用
ExecutorService设置线程池核心数=CPU核心数*2,避免上下文切换开销。 - 批量事务优化:每5000条提交一次事务,Spring中可通过
PlatformTransactionManager手动控制,或在Service层将@Transactional放在批量方法外。 - 连接超时配置:
HikariCP的connectionTimeout设为30000ms,maxLifetime设为1800000ms,防止数据库连接被网络设备切断。
异常处理优先级
- 源库读取失败 → 记录游标位置,下一次重新拉取
- 目标库写入失败 → 回滚当前批次,逐条重试
- 内存溢出 → 启用
XX:+UseG1GC,并降低batch size
常见问题问答(FAQ)
Q1:数据量超过千万时,内存为何仍然溢出?
A:即使使用Cursor,MyBatis默认仍会将数据填充到DefaultResultSetHandler中,需在配置中额外设置:mybatis.configuration.default-fetch-size=-1,强制使用游标模式,同时确保ResultSet类型为TYPE_FORWARD_ONLY。
Q2:迁移过程中业务系统仍在写入MySQL,如何处理增量?
A:本案例专注全量迁移,全量完成后,可通过canal或debezium监听binlog,将增量变更放入Kafka,由Java消费者实时写入PostgreSQL,此部分需做幂等性设计(如使用业务主键或版本号)。
Q3:MySQL和PostgreSQL的JSON字段如何映射?
A:PostgreSQL原生支持JSONB类型,MyBatis可通过typeHandler实现:
@MappedTypes(JSONObject.class)
public class JsonTypeHandler extends BaseTypeHandler<JSONObject> {
@Override
public void setNonNullParameter(PreparedStatement ps, int i, JSONObject parameter, JdbcType jdbcType){
ps.setObject(i, parameter.toJSONString(), Types.OTHER);
}
}
Q4:迁移效率太低,能否压榨硬件性能?
A:可考虑以下路径:①源库侧使用mysqlpump导出压缩文件,目标库使用pg_bulkload加载(跳过WAL日志);②Java端配合CompletableFuture实现生产者-消费者模式,将数据读取和写入解耦,但请注意:过度并行可能导致数据库连接池被打满。
数据迁移的本质是可靠性优先于速度,本文的Java案例展示了如何通过游标查询、批量事务和自定义重试机制,实现从MySQL到PostgreSQL的平滑迁移,实际生产中,还需考虑校验工具(如数据对比脚本)和回滚方案。迁移脚本中的每个try-catch都意味着一次潜在的灾难救援,以上方案已在处理4亿行级别数据时稳定运行,建议结合自身数据库特性微调batch size与并行数。