从MySQL到PostgreSQL:一个Java JDBC程序搞定异构数据库迁移(附完整代码与避坑指南)
从MySQL到PostgreSQL:企业级异构数据库迁移实战指南
当业务系统从单体架构向分布式演进时,数据库异构成为常态。最近接手的一个电商平台改造项目就面临这样的挑战:订单中心采用MySQL,而新建的仓储系统选择了PostgreSQL。两种数据库在事务隔离、锁机制、数据类型等核心特性上的差异,让数据同步成为棘手问题。
1. 迁移方案设计与技术选型
企业级数据迁移绝非简单的SELECT * FROM加INSERT INTO。我们评估了三种主流方案:
- ETL工具:如Kettle、Talend等可视化工具,适合非技术人员但灵活性差
- CDC(变更数据捕获):Debezium等基于日志的方案,适合实时同步但对系统侵入性强
- 定制JDBC程序:完全可控,能处理复杂业务逻辑,本文选择的核心方案
关键决策矩阵:
| 评估维度 | ETL工具 | CDC方案 | JDBC程序 |
|---|---|---|---|
| 开发效率 | ★★★★★ | ★★★☆☆ | ★★☆☆☆ |
| 运行性能 | ★★☆☆☆ | ★★★★☆ | ★★★★★ |
| 业务适配性 | ★★☆☆☆ | ★★★☆☆ | ★★★★★ |
| 运维复杂度 | ★★★☆☆ | ★★☆☆☆ | ★★★★☆ |
// 基础连接示例 - 生产环境务必使用连接池 public class DualConnector { private static final String MYSQL_URL = "jdbc:mysql://mysql-prod:3306/order_db"; private static final String PG_URL = "jdbc:postgresql://pg-warehouse:5432/inventory_db"; public Connection[] getConnections() throws SQLException { Connection[] conns = new Connection[2]; conns[0] = DriverManager.getConnection(MYSQL_URL, "app_user", "加密的密码"); conns[1] = DriverManager.getConnection(PG_URL, "warehouse_user", "加密的密码"); return conns; } }重要提示:生产环境必须配置连接池参数(maxPoolSize、connectionTimeout等),直接使用DriverManager.getConnection会导致性能灾难
2. 数据类型映射的深水区
异构数据库迁移最隐蔽的坑莫过于数据类型差异。上周我们团队就因TIMESTAMP处理不当导致促销活动时间全部错乱。以下是关键映射对照:
数值类型:
- MySQL的
DECIMAL(10,2)→ PostgreSQL的NUMERIC(10,2) - MySQL
INT(11)自增 → PostgreSQLSERIAL
字符串类型:
- MySQL
VARCHAR(255)字符集问题 → PostgreSQLTEXT无长度限制 - MySQL的
utf8mb4才是真正的UTF-8
日期时间:
- MySQL
DATETIME无时区 → PostgreSQLTIMESTAMP WITH TIME ZONE - MySQL
ON UPDATE CURRENT_TIMESTAMP语法在PG中完全不同
-- PostgreSQL需要特殊处理自增ID CREATE TABLE products ( id SERIAL PRIMARY KEY, -- 替代MySQL的AUTO_INCREMENT name VARCHAR(100) NOT NULL, price NUMERIC(10,2) CHECK (price > 0) );3. 高性能批量迁移实战
当需要迁移百万级数据时,逐条插入会导致迁移时间呈指数增长。我们通过三种优化手段将迁移速度提升37倍:
- 批处理操作:利用
addBatch()和executeBatch() - 事务分片:每1万条提交一次,避免超大事务
- 并行迁移:按时间范围切分数据并行处理
// 优化后的批量插入代码片段 public void batchInsert(List<Product> products, Connection pgConn) throws SQLException { final int BATCH_SIZE = 1000; String sql = "INSERT INTO products (name, price, stock) VALUES (?, ?, ?)"; try (PreparedStatement pstmt = pgConn.prepareStatement(sql)) { for (int i = 0; i < products.size(); i++) { Product p = products.get(i); pstmt.setString(1, p.getName()); pstmt.setBigDecimal(2, p.getPrice()); pstmt.setInt(3, p.getStock()); pstmt.addBatch(); if (i % BATCH_SIZE == 0 || i == products.size() - 1) { pstmt.executeBatch(); pgConn.commit(); // 分批次提交 } } } }性能对比测试(迁移10万条商品数据):
| 方案 | 耗时(ms) | 内存峰值(MB) |
|---|---|---|
| 单条插入 | 182,456 | 1,024 |
| 纯批处理 | 23,781 | 512 |
| 批处理+事务分片 | 4,932 | 256 |
4. 生产环境必须的增强特性
基础迁移代码只能应付Demo,要上线还需要以下企业级功能:
健壮性保障:
- 断点续传:记录最后成功ID,程序重启后继续
- 数据校验:CRC32校验和比对源库与目标库
- 异常处理:网络闪断重试机制
可观测性:
- 埋点监控:迁移速率、数据差异等指标
- 详细日志:记录跳过或失败的记录详情
- 预警机制:超过阈值自动告警
// 断点续传实现示例 public class MigrationState { private static final String STATE_FILE = "/data/migration.state"; public void saveLastId(long lastId) throws IOException { Files.write(Paths.get(STATE_FILE), String.valueOf(lastId).getBytes()); } public long loadLastId() throws IOException { if (Files.exists(Paths.get(STATE_FILE))) { String id = new String(Files.readAllBytes(Paths.get(STATE_FILE))); return Long.parseLong(id.trim()); } return 0L; // 首次运行从0开始 } }经验之谈:实际项目中我们增加了Redis分布式锁,防止多个迁移实例同时运行导致数据重复
5. 进阶:双向同步解决方案
当业务需要MySQL和PostgreSQL保持实时双向同步时,单纯的迁移程序就不够用了。我们最终采用的架构:
- 变更捕获层:MySQL用binlog,PG用逻辑解码
- 消息队列缓冲:Kafka作为中间件解耦
- 冲突解决策略:时间戳+业务规则判断最后更新
# 简化的冲突解决伪代码 def resolve_conflict(mysql_row, pg_row): mysql_time = mysql_row['updated_at'] pg_time = pg_row['updated_at'] if mysql_time > pg_time: return mysql_row elif pg_time > mysql_time: return pg_row else: # 按业务优先级处理 if mysql_row['version'] > pg_row['version']: return mysql_row else: return pg_row这套方案最终支撑了日均2000万次的跨库数据同步,延迟控制在500ms以内。关键点在于合理设置批量处理大小和消费者线程数——太大导致延迟增加,太小则浪费资源。
