数据库迁移策略:零停机Schema变更与数据同步实战

全面讲解数据库迁移的核心策略与最佳实践,涵盖在线Schema变更、大表迁移、跨数据库迁移、数据同步等场景,提供gh-ost、pg_repack、Debezium等工具的实战案例。

引言

数据库迁移是后端开发中最具挑战性的任务之一。不当的迁移策略可能导致服务停机、数据丢失或性能问题。

本文将介绍零停机数据库迁移的核心策略和实战工具。

迁移策略分类

按停机时间分类

┌─────────────────────────────────────────────────────────┐
│                   迁移策略对比                           │
│                                                         │
│  1. 停机迁移(Offline Migration)                       │
│     - 停止服务 → 迁移 → 重启服务                        │
│     - 适用:非关键系统、维护窗口充足                      │
│     - 缺点:服务不可用                                   │
│                                                         │
│  2. 在线迁移(Online Migration)                        │
│     - 服务持续运行,后台迁移                             │
│     - 适用:24/7在线服务                                 │
│     - 工具:gh-ost、pt-online-schema-change             │
│                                                         │
│  3. 双写迁移(Dual-Write Migration)                    │
│     - 同时写入新旧库,逐步切换                           │
│     - 适用:跨数据库迁移                                 │
│     - 复杂度:高                                         │
└─────────────────────────────────────────────────────────┘

MySQL在线DDL

gh-ost:GitHub开源的在线Schema变更工具

# 安装gh-ost
brew install gh-ost  # macOS
# 或从GitHub下载:https://github.com/github/gh-ost/releases

# 添加列(在线执行)
gh-ost \
  --host="db-host" \
  --port=3306 \
  --user="migration_user" \
  --password="password" \
  --database="mydb" \
  --table="orders" \
  --alter="ADD COLUMN payment_method VARCHAR(50) DEFAULT 'credit_card'" \
  --allow-on-master \
  --execute

# 修改列类型
gh-ost \
  --host="db-host" \
  --port=3306 \
  --user="migration_user" \
  --password="password" \
  --database="mydb" \
  --table="users" \
  --alter="MODIFY COLUMN phone_number VARCHAR(20)" \
  --allow-on-master \
  --execute

# 添加索引
gh-ost \
  --host="db-host" \
  --port=3306 \
  --user="migration_user" \
  --password="password" \
  --database="mydb" \
  --table="orders" \
  --alter="ADD INDEX idx_user_created (user_id, created_at)" \
  --allow-on-master \
  --execute

# 删除列(危险操作,需要确认)
gh-ost \
  --host="db-host" \
  --port=3306 \
  --user="migration_user" \
  --password="password" \
  --database="mydb" \
  --table="orders" \
  --alter="DROP COLUMN legacy_field" \
  --allow-on-master \
  --approve-renamed-columns \
  --execute
# gh-ost高级选项
gh-ost \
  --host="db-host" \
  --port=3306 \
  --user="migration_user" \
  --password="password" \
  --database="mydb" \
  --table="orders" \
  --alter="ADD COLUMN new_field INT" \
  --allow-on-master \
  --execute \
  --chunk-size=1000 \                    # 每批处理的行数
  --max-lag-millis=1500 \                # 最大复制延迟(毫秒)
  --throttle-query="SELECT timestampdiff(SECOND, MIN(last_update), NOW()) FROM heartbeat.heartbeat" \
  --throttle-control-replicas="replica1,replica2" \
  --cut-over=default \                   # 切换策略
  --panic-flag-file=/tmp/gh-ost.panic \  # 紧急停止标志
  --postpone-cut-over-flag-file=/tmp/gh-ost.postpone  # 延迟切换

pt-online-schema-change(Percona Toolkit)

# 安装Percona Toolkit
brew install percona-toolkit  # macOS
apt-get install percona-toolkit  # Ubuntu

# 添加列
pt-online-schema-change \
  --alter "ADD COLUMN status VARCHAR(20) DEFAULT 'pending'" \
  D=mydb,t=orders \
  --user=migration_user \
  --password=password \
  --execute

# 修改列
pt-online-schema-change \
  --alter "MODIFY COLUMN description TEXT" \
  D=mydb,t=products \
  --user=migration_user \
  --password=password \
  --execute

# 添加外键
pt-online-schema-change \
  --alter "ADD CONSTRAINT fk_user FOREIGN KEY (user_id) REFERENCES users(id)" \
  D=mydb,t=orders \
  --user=migration_user \
  --password=password \
  --execute

PostgreSQL在线迁移

pg_repack:在线重建表和索引

# 安装pg_repack
apt-get install postgresql-14-pgrepack  # Ubuntu
brew install pg_repack  # macOS

# 重建表(消除碎片,回收空间)
pg_repack --table=orders --dbname=mydb

# 重建索引
pg_repack --index=idx_orders_user_id --dbname=mydb

# 重建整个数据库的所有表
pg_repack --dbname=mydb --no-superuser-check

# 仅重建碎片化严重的表
pg_repack --table=large_table --dbname=mydb --jobs=4

在线添加列和索引

-- PostgreSQL 11+ 支持在线添加带默认值的列
ALTER TABLE orders ADD COLUMN status VARCHAR(20) DEFAULT 'pending';

-- 并发创建索引(不阻塞读写)
CREATE INDEX CONCURRENTLY idx_orders_status ON orders(status);

-- 并发删除索引
DROP INDEX CONCURRENTLY idx_orders_old;

-- 在线修改列类型(使用USING)
ALTER TABLE users 
  ALTER COLUMN age TYPE INTEGER 
  USING age::INTEGER;

大表分区

-- 创建分区表
CREATE TABLE orders (
    id BIGSERIAL,
    user_id BIGINT,
    created_at TIMESTAMP,
    total_amount DECIMAL(10, 2)
) PARTITION BY RANGE (created_at);

-- 创建月度分区
CREATE TABLE orders_2026_01 PARTITION OF orders
    FOR VALUES FROM ('2026-01-01') TO ('2026-02-01');

CREATE TABLE orders_2026_02 PARTITION OF orders
    FOR VALUES FROM ('2026-02-01') TO ('2026-03-01');

-- 自动创建分区(使用pg_partman)
CREATE EXTENSION pg_partman;

SELECT partman.create_parent(
    p_parent_table := 'public.orders',
    p_control := 'created_at',
    p_type := 'native',
    p_interval := '1 month',
    p_premake := 3
);

跨数据库迁移

双写策略

// 双写迁移实现
type DualWriteRepository struct {
    oldDB *gorm.DB
    newDB *gorm.DB
    
    // 迁移状态
    migrationEnabled bool
    readFromNew      bool
}

func (r *DualWriteRepository) CreateOrder(ctx context.Context, order *Order) error {
    // 1. 写入旧库(主库)
    if err := r.oldDB.WithContext(ctx).Create(order).Error; err != nil {
        return err
    }
    
    // 2. 如果启用迁移,同时写入新库
    if r.migrationEnabled {
        // 异步写入新库,避免影响主流程
        go func() {
            if err := r.newDB.WithContext(ctx).Create(order).Error; err != nil {
                log.Errorf("Failed to write to new DB: %v", err)
                // 记录失败,后续可以通过补偿任务重试
                recordMigrationFailure(order.ID, err)
            }
        }()
    }
    
    return nil
}

func (r *DualWriteRepository) GetOrder(ctx context.Context, id string) (*Order, error) {
    // 根据配置决定从哪个库读取
    if r.readFromNew {
        return r.getOrderFromNewDB(ctx, id)
    }
    return r.getOrderFromOldDB(ctx, id)
}

// 迁移步骤
func migrateToNewDB() {
    // 阶段1:启用双写,读旧库
    repo.migrationEnabled = true
    repo.readFromNew = false
    
    // 阶段2:历史数据迁移(使用CDC或批处理)
    migrateHistoricalData()
    
    // 阶段3:验证数据一致性
    validateDataConsistency()
    
    // 阶段4:切换读流量到新库
    repo.readFromNew = true
    
    // 阶段5:观察一段时间,确认无问题
    time.Sleep(7 * 24 * time.Hour)
    
    // 阶段6:停止双写,完全切换到新库
    repo.migrationEnabled = false
}

CDC(Change Data Capture)数据同步

# Debezium MySQL连接器配置
{
  "name": "orders-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "database.hostname": "mysql-host",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "password",
    "database.server.id": "184054",
    "database.server.name": "mydb",
    "database.include.list": "mydb",
    "table.include.list": "mydb.orders",
    
    "topic.prefix": "cdc",
    "schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
    "schema.history.internal.kafka.topic": "schema-changes.mydb",
    
    "include.schema.changes": "true",
    "snapshot.mode": "initial"
  }
}
// 消费CDC事件,同步到新库
type CDCConsumer struct {
    kafkaReader *kafka.Reader
    newDB       *gorm.DB
}

func (c *CDCConsumer) Consume(ctx context.Context) error {
    for {
        msg, err := c.kafkaReader.ReadMessage(ctx)
        if err != nil {
            return err
        }
        
        var event CDCEvent
        if err := json.Unmarshal(msg.Value, &event); err != nil {
            log.Errorf("Failed to unmarshal CDC event: %v", err)
            continue
        }
        
        switch event.Operation {
        case "c": // Create
            c.handleCreate(event)
        case "u": // Update
            c.handleUpdate(event)
        case "d": // Delete
            c.handleDelete(event)
        }
    }
}

func (c *CDCConsumer) handleCreate(event CDCEvent) {
    var order Order
    json.Unmarshal(event.After, &order)
    
    c.newDB.Create(&order)
}

func (c *CDCConsumer) handleUpdate(event CDCEvent) {
    var order Order
    json.Unmarshal(event.After, &order)
    
    c.newDB.Model(&Order{}).Where("id = ?", order.ID).Updates(&order)
}

func (c *CDCConsumer) handleDelete(event CDCEvent) {
    var order Order
    json.Unmarshal(event.Before, &order)
    
    c.newDB.Delete(&Order{}, "id = ?", order.ID)
}

大表迁移策略

分批迁移

// 大表分批迁移
type BatchMigrator struct {
    oldDB *gorm.DB
    newDB *gorm.DB
    batchSize int
}

func (m *BatchMigrator) MigrateTable(ctx context.Context) error {
    var lastID int64 = 0
    
    for {
        // 从旧库分批读取数据
        var orders []Order
        if err := m.oldDB.WithContext(ctx).
            Where("id > ?", lastID).
            Order("id").
            Limit(m.batchSize).
            Find(&orders).Error; err != nil {
            return err
        }
        
        // 没有更多数据,迁移完成
        if len(orders) == 0 {
            break
        }
        
        // 批量插入到新库
        if err := m.newDB.WithContext(ctx).
            CreateInBatches(orders, m.batchSize).Error; err != nil {
            return err
        }
        
        // 更新lastID
        lastID = orders[len(orders)-1].ID
        
        log.Printf("Migrated %d records, lastID: %d", len(orders), lastID)
        
        // 避免对数据库造成过大压力
        time.Sleep(100 * time.Millisecond)
    }
    
    return nil
}

数据一致性验证

// 数据一致性验证
type DataValidator struct {
    oldDB *gorm.DB
    newDB *gorm.DB
}

func (v *DataValidator) Validate(ctx context.Context) error {
    // 1. 验证总数
    var oldCount, newCount int64
    v.oldDB.Model(&Order{}).Count(&oldCount)
    v.newDB.Model(&Order{}).Count(&newCount)
    
    if oldCount != newCount {
        return fmt.Errorf("count mismatch: old=%d, new=%d", oldCount, newCount)
    }
    
    // 2. 抽样验证数据
    var sampleIDs []int64
    v.oldDB.Model(&Order{}).
        Select("id").
        Order("RANDOM()").
        Limit(1000).
        Pluck("id", &sampleIDs)
    
    for _, id := range sampleIDs {
        var oldOrder, newOrder Order
        v.oldDB.First(&oldOrder, id)
        v.newDB.First(&newOrder, id)
        
        if !reflect.DeepEqual(oldOrder, newOrder) {
            log.Errorf("Data mismatch for order %d", id)
        }
    }
    
    // 3. 验证关键指标
    var oldTotal, newTotal float64
    v.oldDB.Model(&Order{}).Select("SUM(total_amount)").Scan(&oldTotal)
    v.newDB.Model(&Order{}).Select("SUM(total_amount)").Scan(&newTotal)
    
    if math.Abs(oldTotal-newTotal) > 0.01 {
        return fmt.Errorf("total amount mismatch: old=%f, new=%f", oldTotal, newTotal)
    }
    
    log.Info("Data validation passed")
    return nil
}

迁移回滚策略

回滚计划

// 迁移回滚策略
type MigrationRollback struct {
    backupDB *gorm.DB
    mainDB   *gorm.DB
}

// 迁移前创建备份点
func (r *MigrationRollback) CreateBackupPoint(ctx context.Context) error {
    // 1. 记录当前时间点
    r.backupTimestamp = time.Now()
    
    // 2. 创建数据库快照(使用云服务商的快照功能)
    // AWS RDS: aws rds create-db-snapshot
    // GCP Cloud SQL: gcloud sql backups create
    
    // 3. 开启binlog记录(MySQL)
    // 确保可以通过binlog回滚到特定时间点
    
    return nil
}

// 回滚到备份点
func (r *MigrationRollback) Rollback(ctx context.Context) error {
    // 1. 停止应用写入
    // 2. 从快照恢复数据库
    // 3. 重放binlog到迁移前的时间点
    // 4. 验证数据完整性
    // 5. 重启应用
    
    return nil
}

Schema 变更版本管理

生产环境中的 Schema 变更需要版本化管理,避免多人同时修改导致冲突:

迁移脚本命名规范:
  V001__create_users_table.sql
  V002__add_user_index.sql
  V003__add_order_table.sql

工具对比:
┌─────────────┬─────────┬──────────┬──────────┐
│ 工具        │ 语言    │ 回滚支持 │ 校验和   │
├─────────────┼─────────┼──────────┼──────────┤
│ Flyway      │ SQL     │ ✅       │ ✅       │
│ Liquibase   │ XML/YAML│ ✅       │ ✅       │
│ golang-migrate│ Go   │ ✅       │ ❌       │
│ Atlas       │ HCL     │ ✅       │ ✅       │
└─────────────┴─────────┴──────────┴──────────┘

Flyway 的 repeatable 迁移适合视图和存储过程更新,versioned 迁移负责 Schema 变更。flyway validate 在部署前检查脚本完整性,防止生产环境被篡改。

迁移风险与回滚实战

常见陷阱

陷阱场景后果预防
大表加锁ALTER TABLE 无 ONLINE 选项读写阻塞数小时使用 gh-ost/pt-osc
索引名冲突不同环境索引命名不一致迁移失败统一命名规范
字符集变更从 latin1 改为 utf8mb4数据截断先验证字段长度
默认值陷阱ALTER 添加 NOT NULL 无默认值现有行填充失败先填 NULL,再改 NOT NULL

回滚决策矩阵

发现问题时的判断:

数据一致性受损?
  → 是:立即回滚,停止写入,从备份恢复
  → 否:继续观察

性能劣化超过 50%?
  → 是:回滚 Schema 变更(如新增索引可快速 DROP)
  → 否:可接受,等低峰期调整

业务功能异常?
  → 是:优先代码回滚,Schema 变更通常不可逆需谨慎
  → 否:继续监控

Schema 变更回滚通常比代码回滚复杂。添加列可以删除,但删除列后无法「加回来」原有数据。

总结

数据库迁移需要谨慎规划和执行:

  1. 在线DDL工具:gh-ost(MySQL)、pg_repack(PostgreSQL)实现零停机Schema变更
  2. 跨库迁移:双写策略 + CDC数据同步,逐步切换流量
  3. 大表迁移:分批迁移 + 数据一致性验证
  4. 回滚计划:迁移前创建备份,准备回滚方案

关键原则:小步快跑、充分验证、随时可回滚。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「backend」更多文章

  1. 零信任安全架构:从边界防御到身份中心的安全范式
  2. 混沌工程实践:构建高可用系统的故障注入与弹性测试
  3. 流式数据处理:Kafka Streams与Flink实战指南