不停机数据迁移

2023-02-10T14:11:02+08:00 | 15分钟阅读 | 更新于 2026-02-10T14:11:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 解释为什么微服务化后必须迁移数据,说清不停机迁移的核心难点。
  2. 复述不停机迁移的四阶段模型:SRC_ONLY → SRC_FIRST → DST_FIRST → DST_ONLY。
  3. 设计数据导出导入、全量校验、增量校验、数据修复的完整方案。
  4. 理解基于 GORM ConnPool 的双写实现,能写出生产可用的双写代码。
  5. 处理迁移中的并发问题、性能优化与回滚方案,并能讲成一段面试故事。

前置知识(如果下面任意一点生疏,先回看对应章):

  • 第02章 Gin + GORM:知道事务、First/Save/批量操作怎么写。
  • 第07章 Kafka:知道消息生产者/消费者、削峰与解耦(修复环节会用到)。
  • 第10章 单体拆分微服务:知道「为什么要把数据库也拆开」。
  • 基本的乐观锁(version 字段)、UPSERT 概念。

本章你会动手做的事

  • mysqldump 把一张表导出来、再导入到一个新库,跑通「初始化目标表」。
  • 写一个 DoubleWritePool,手动切 PatternSRCFirstPatternDSTFirstPatternDSTOnly,观察读写打向哪个库。
  • 触发一次全量校验,故意改一条目标表数据,看修复流程怎么把它改回来。

一、核心概念讲解

1.1 为什么要数据迁移

微服务化之后,代码层面已经拆分为独立服务,但数据库还在共享(逻辑上和物理上都是同一个)。微服务理论有一个广泛认同的观点:

某个微服务的数据应该单独存储。

因此,正常在微服务化之后,还要把数据分离、独立存储。

数据迁移是几乎每个后端工程师职业生涯都会遇到的任务,非常适合作为面试亮点。

1.2 停机迁移 vs 不停机迁移

停机迁移最自然:停下应用 → 从老库导出数据 → 导入新库 → 启动应用。

  • 优点:方案简单;
  • 缺点:应用停下来了!如果数据量大,导出和导入都很慢,应用要停很久,严重影响用户体验。

不停机迁移是核心应用必须采用的方案,难点在于:

  • 迁移过程中数据始终在变动(业务一边写一边迁);
  • 不能对数据库造成太大压力,否则影响线上业务。

哲学类比——蜗牛悖论:你每次往前追一段,蜗牛也往前走一段,永远追不上。不停机迁移的核心难点就在这里:不管你怎么迁移,t1 时刻源库和目标库一致了,但应用立刻又修改了源库,目标库又落后了。

sequenceDiagram
    participant App as 业务应用
    participant SRC as 源库
    participant DST as 目标库
    App->>SRC: 写入数据
    Note over SRC,DST: t1 时刻两库一致
    App->>SRC: 又改了源库
    Note over DST: 目标库又落后了

1.3 异构数据迁移

如果源表和目标表结构不一样(甚至数据库都不一样,如 MySQL → MongoDB),叫异构数据迁移。难点:

  • 数据转换:类型不兼容、精度损失、长度溢出;
  • 完整性校验:外键约束、业务关系难以验证。

异构数据迁移通常只能定制校验逻辑,无法用通用方案。

1.4 不停机迁移的四阶段模型

阶段名称写入读取数据基准
第一阶段SRC_ONLY源表源表-
第二阶段SRC_FIRST双写源表 + 目标表源表源表为准
第三阶段DST_FIRST双写源表 + 目标表目标表目标库为准
第四阶段DST_ONLY目标表目标表-

关键设计:第二、三阶段都是双写,但读哪个表、修复时以哪个表为准会切换。这样可以在出问题时随时切回上一阶段。

stateDiagram-v2
    [*] --> SRC_ONLY: 初始
    SRC_ONLY --> SRC_FIRST: 开启双写
    SRC_FIRST --> DST_FIRST: 校验通过,切读
    DST_FIRST --> DST_ONLY: 校验通过,切写
    DST_ONLY --> [*]
    SRC_FIRST --> SRC_ONLY: 出问题回滚
    DST_FIRST --> SRC_FIRST: 出问题回滚

1.5 具体迁移步骤

  1. 创建目标表;
  2. 用源表数据初始化目标表;
  3. 执行一次校验,并用源表数据修复目标表(可选,开启双写前减少差异);
  4. 业务代码开启双写,读源表,先写源表,源表为准(SRC_FIRST);
  5. 开启增量校验和数据修复,保持一段时间;
  6. 切换双写顺序,读目标表,先写目标表,目标表为准(DST_FIRST);
  7. 继续保持增量校验和修复;
  8. 切换为目标表单写,读写都只操作目标表(DST_ONLY)。

二、数据导出与导入

2.1 使用 mysqldump 导出

mysqldump 是 MySQL 自带的命令行工具,原本用于备份:

mysqldump \
  -h 127.0.0.1 \           # host,必须指定,否则默认走 /tmp/mysql.sock 会报错
  --port 3306 \
  -u root -p \
  webook interactive > intr.sql

参数说明:

  • -h:指定 host;
  • --port:指定端口;
  • --result-file:也可用此参数指定输出文件路径。

2.2 导入到目标库

-- 创建目标数据库
CREATE DATABASE IF NOT EXISTS webook_intr;
# 命令行连上 MySQL 后直接 source
mysql -u root -p webook_intr -e "source /path/to/intr.sql"

2.3 异构数据初始化

异构数据不能直接导入,思路是写一个批量迁移程序:

// webook/internal/migration/init.go (示意)
func MigrateData(ctx context.Context, srcDB, dstDB *gorm.DB) error {
    var srcRecords []SrcRecord
    // 步骤 1:分批读取源库数据
    err := srcDB.WithContext(ctx).Model(&SrcRecord{}).
        FindInBatches(&srcRecords, 500, func(tx *gorm.DB, batch int) error {
            // 步骤 2:转换为目标库格式
            for _, src := range srcRecords {
                dst := convertToDst(src)
                // 步骤 3:写入目标库
                if err := dstDB.Create(&dst).Error; err != nil {
                    return err
                }
            }
            return nil
        }).Error
    return err
}

并发优化策略:

  • 表内拆分:N 个 goroutine,分别负责 id % N == 0/1/.../N-1 的数据;
  • 每个 goroutine 负责不同的表
  • 两种方式可混合使用;
  • 控制 goroutine 数量,避免影响线上应用。

三、数据校验逻辑

3.1 全量校验与修复

初始化目标表后,可选地执行一次全量校验和修复(开启双写前减少两库差异)。

全量校验基本思路:从源表取数据 → 按主键去目标表找对应数据 → 比较所有字段。

坑点

  • 数据库类型能否都转成 Go 类型?
  • 转成 Go 类型后是否可比较?
  • 浮点数精度损失问题?

3.2 方案选型

方案优点缺点
每张表写一个 DAO + 比较方法直观表多时要命
泛型通用查询 + Entity 实现比较通用需要实体定义
[]byte 接收数据直接比较通用类型信息丢失

课程采用方案二。

3.3 校验方案——Validator 定义

// webook/internal/migration/validator.go
package migration

import "gorm.io/gorm"

// Entity 是校验对象的接口
type Entity interface {
    // ID 主键,整个校验体系围绕 ID 进行
    ID() int64
    // TableName 表名,修复时需要
    TableName() string
    // CompareTo 比较两个实体是否一致
    // 由实现者负责,明确知道每个字段怎么比较、要不要忽略某些列
    CompareTo(other Entity) error
    // Columns 所有列名,修复时用
    Columns() []string
}

// Validator 泛型校验器
type Validator[T Entity] struct {
    baseDB    *gorm.DB // 基准库(第二阶段=源表,第三阶段=目标表)
    targetDB  *gorm.DB // 对比库
    direction Direction // SRC 或 DST
}

// Direction 数据方向
type Direction uint8

const (
    DirectionSRC Direction = iota // 以源表为准
    DirectionDST                  // 以目标表为准
)

3.4 校验逻辑实现

// Validate 执行校验,offset 是从哪个 ID 开始
func (v *Validator[T]) Validate(ctx context.Context, offset int64) (int64, error) {
    var base T
    // 步骤 1:从 base 中取一条数据
    err := v.baseDB.WithContext(ctx).
        Where("id > ?", offset).
        First(&base).Error
    if err == gorm.ErrRecordNotFound {
        // 校验完毕
        return offset, nil
    }
    if err != nil {
        return offset, err
    }

    // 步骤 2:从 target 中找对应数据
    var target T
    err = v.targetDB.WithContext(ctx).
        Where("id = ?", base.ID()).First(&target).Error
    if err == gorm.ErrRecordNotFound {
        // target 中没有,需要修复(insert)
        v.notifyFix(base.ID(), "target_missing")
        return base.ID(), nil
    }
    if err != nil {
        // 查询异常,记录日志,继续下一条
        log.Warn("查询 target 出错", err)
        return base.ID(), nil
    }

    // 步骤 3:比较数据
    if cmpErr := base.CompareTo(target); cmpErr != nil {
        // 数据不一致,需要修复(update)
        v.notifyFix(base.ID(), "neq")
    }
    return base.ID(), nil
}

3.5 反向校验——target 中多出的数据

在 SRC_FIRST 阶段,可能源库做了硬删除(DELETE),导致目标库多出数据(软删除是 UPDATE,不会有这个问题)。需要反向校验:

// ReverseValidate 反向校验:找 target 中有但 base 中没有的记录
func (v *Validator[T]) ReverseValidate(ctx context.Context) error {
    var ids []int64
    // 批量取 target 中的 ID
    cursor := int64(0)
    for {
        ids = ids[:0]
        err := v.targetDB.WithContext(ctx).
            Where("id > ?", cursor).
            Limit(500).Pluck("id", &ids).Error
        if err != nil || len(ids) == 0 {
            return err
        }

        // 在 base 中找这些 ID
        var found []int64
        v.baseDB.WithContext(ctx).
            Where("id IN ?", ids).Pluck("id", &found).Error
        foundSet := make(map[int64]struct{}, len(found))
        for _, id := range found {
            foundSet[id] = struct{}{}
        }

        // 找出 target 有但 base 没有的,需要修复(delete from target)
        for _, id := range ids {
            if _, ok := foundSet[id]; !ok {
                v.notifyFix(id, "base_missing")
            }
        }
        cursor = ids[len(ids)-1]
    }
}

3.6 异构数据的校验

异构数据库校验无法用通用方案:

  • 源表查出数据后无法用主键去目标表找;
  • 可考虑用唯一索引、外键查找;
  • 字段可能被丢弃或重新计算;
  • 只能定制校验逻辑。

3.7 校验调度时机

  • 业务低峰期运行:定时任务,设置在低峰期;
  • 动态判定负载:负载高挂起,负载低继续。

四、数据修复

4.1 修复消息结构

校验发现不一致时,发送 Kafka 消息:

type FixMsg struct {
    ID        int64     `json:"id"`        // 主键
    Type      FixType   `json:"type"`      // 不一致类型:target_missing / neq / base_missing
    Direction Direction `json:"direction"` // 以 SRC 为准还是 DST 为准
    Table     string    `json:"table"`     // 表名
}

type FixType uint8

const (
    FixTypeTargetMissing FixType = iota // 目标表缺数据,执行 insert
    FixTypeNeq                          // 数据不相等,执行 update
    FixTypeBaseMissing                  // 目标表多数据,执行 delete
)

4.2 修复逻辑

// Fix 修复一条数据
func (f *Fixer[T]) Fix(ctx context.Context, msg FixMsg) error {
    var base T
    // 再次查找 base 中的最新数据(避免使用校验时的旧数据)
    err := f.baseDB.WithContext(ctx).
        Where("id = ?", msg.ID).First(&base).Error
    if err == gorm.ErrRecordNotFound {
        // base 中没有了,删除 target 中的数据
        return f.targetDB.WithContext(ctx).
            Where("id = ?", msg.ID).Delete(base).Error
    }
    if err != nil {
        return err
    }
    // base 中有,执行 UPSERT(存在则更新,不存在则插入)
    return f.targetDB.WithContext(ctx).Save(&base).Error
}
flowchart TD
    V[校验发现不一致] --> K[发 Kafka FixMsg]
    K --> C[消费者取 base 最新数据]
    C --> D{base 有这条?}
    D -- 有 --> U[UPSERT 到 target]
    D -- 无 --> X[DELETE target 中的记录]

4.3 简化版——不用区分类型

实际上,修复时不需要区分校验出来的不一致类型是 target_missing 还是 neq

  • 直接去 base 找数据 → 找到就 UPSERT,没找到就 DELETE。
// SimplifiedFix 简化版修复
func (f *Fixer[T]) SimplifiedFix(ctx context.Context, id int64) error {
    var base T
    err := f.baseDB.WithContext(ctx).Where("id = ?", id).First(&base).Error
    if err == gorm.ErrRecordNotFound {
        return f.targetDB.WithContext(ctx).Where("id = ?", id).Delete(base).Error
    }
    if err != nil {
        return err
    }
    return f.targetDB.WithContext(ctx).Save(&base).Error
}

4.4 为什么用 UPSERT?并发问题

核心原因:双写阶段,校验时 target 没有这条数据,但紧接着你修复时,业务的双写已经写进去了。

所以必须用 UPSERT(INSERT … ON DUPLICATE KEY UPDATE 或 GORM 的 Save),避免唯一键冲突。

⚠️ 新手必踩的坑:修复用普通 INSERT 而不是 UPSERT。校验的那一刻 target 没数据,可等你发修复消息、消费者去修的时候,业务双写可能已经把这条写进 target 了。这时再 INSERT 就会主键/唯一键冲突报错。用 Save(UPSERT 语义)就能「有则更新、无则插入」,天然规避。

4.5 为什么用 Kafka 解耦

修复阶段用 Kafka 主要为了保护目标表

  • 通过 Kafka 削峰,控制消费者速率,间接控制目标表的写入速率;
  • 易于横向扩展消费者数量;
  • 面试时也"更高级"。

五、双写实现——基于 GORM ConnPool

5.1 三种双写思路

思路优缺点
直接改 DAO侵入式,易引入 BUG,易遗漏
GORM Hook/CallbackHook 要一个个结构体定义;Callback 收到的 DB 固定,无法动态切换
GORM ConnPool高级但优雅,能动态切换 pattern,能控制事务

5.2 ConnPool 接口

// gorm.io/gorm 中的底层接口
type ConnPool interface {
    PrepareContext(ctx context.Context, query string) (*sql.Stmt, error)
    ExecContext(ctx context.Context, query string, args ...interface{}) (sql.Result, error)
    QueryContext(ctx context.Context, query string, args ...interface{}) (*sql.Rows, error)
    QueryRowContext(ctx context.Context, query string, args ...interface{}) *sql.Row
}

它是 GORM 直接和数据库打交道的出口。

5.3 DoubleWritePool 实现

// webook/internal/migration/double_write_pool.go
package migration

import (
    "context"
    "database/sql"
    "sync/atomic"
)

// Pattern 双写模式
type Pattern uint8

const (
    PatternSRCOnly   Pattern = iota // 只读写源库
    PatternSRCFirst                 // 双写,读源库,源库为准
    PatternDSTFirst                 // 双写,读目标库,目标库为准
    PatternDSTOnly                  // 只读写目标库
)

// DoubleWritePool 双写连接池
type DoubleWritePool struct {
    src     *gorm.DB // 源库
    dst     *gorm.DB // 目标库
    pattern atomic.Int32
}

func NewDoubleWritePool(src, dst *gorm.DB) *DoubleWritePool {
    p := &DoubleWritePool{src: src, dst: dst}
    p.pattern.Store(int32(PatternSRCOnly))
    return p
}

func (p *DoubleWritePool) UpdatePattern(pat Pattern) {
    p.pattern.Store(int32(pat))
}

// ExecContext 处理增删改
func (p *DoubleWritePool) ExecContext(ctx context.Context, query string, args ...interface{}) (sql.Result, error) {
    pat := Pattern(p.pattern.Load())
    switch pat {
    case PatternSRCOnly:
        return p.src.ExecContext(ctx, query, args...)
    case PatternDSTOnly:
        return p.dst.ExecContext(ctx, query, args...)
    case PatternSRCFirst, PatternDSTFirst:
        // 双写:先写主库,再写从库
        res, err := p.src.ExecContext(ctx, query, args...)
        if err != nil {
            return res, err
        }
        // 写 dst 失败不返回错误,等校验修复程序兜底
        if _, e := p.dst.ExecContext(ctx, query, args...); e != nil {
            log.Warn("双写 dst 失败", e)
        }
        return res, nil
    }
    return nil, fmt.Errorf("unknown pattern: %d", pat)
}

// QueryContext 处理查询
func (p *DoubleWritePool) QueryContext(ctx context.Context, query string, args ...interface{}) (*sql.Rows, error) {
    pat := Pattern(p.pattern.Load())
    switch pat {
    case PatternSRCOnly, PatternSRCFirst:
        return p.src.QueryContext(ctx, query, args...)
    case PatternDSTOnly, PatternDSTFirst:
        return p.dst.QueryContext(ctx, query, args...)
    }
    return nil, fmt.Errorf("unknown pattern: %d", pat)
}

// QueryRowContext 单行查询
func (p *DoubleWritePool) QueryRowContext(ctx context.Context, query string, args ...interface{}) *sql.Row {
    pat := Pattern(p.pattern.Load())
    switch pat {
    case PatternSRCOnly, PatternSRCFirst:
        return p.src.QueryRowContext(ctx, query, args...)
    case PatternDSTOnly, PatternDSTFirst:
        return p.dst.QueryRowContext(ctx, query, args...)
    }
    return p.src.QueryRowContext(ctx, query, args...)
}

⚠️ 新手必踩的坑:双写「先写源、dst 失败吞错误」会不会丢数据?不会永久丢。dst 写失败只是这次没写进去,后面有「全量/增量校验 + 修复」兜底把它补回来。这正是双写方案敢「dst 失败不返回 error」的底气——用「最终一致」换「主流程绝对不被拖垮」。

5.4 双写事务

业务经常用事务,双写时意味着源库和目标库各开一个事务:

// DoubleWriteTx 双写事务
type DoubleWriteTx struct {
    srcTx *gorm.DB
    dstTx *gorm.DB
    pattern Pattern
}

func (t *DoubleWriteTx) Commit() error {
    if err := t.srcTx.Commit().Error; err != nil {
        return err
    }
    if t.pattern == PatternSRCFirst || t.pattern == PatternDSTFirst {
        if err := t.dstTx.Commit().Error; err != nil {
            // dst 提交失败,记录日志,等校验修复兜底
            log.Warn("dst 事务提交失败", err)
        }
    }
    return nil
}

func (t *DoubleWriteTx) Rollback() error {
    if err := t.srcTx.Rollback().Error; err != nil {
        return err
    }
    if t.pattern == PatternSRCFirst || t.pattern == PatternDSTFirst {
        if err := t.dstTx.Rollback().Error; err != nil {
            log.Warn("dst 事务回滚失败", err)
        }
    }
    return nil
}

5.5 容错策略(以 SRC_FIRST 为例)

场景处理方式
写 SRC 成功,写 DST 失败直接触发修复,或不管,等后续校验找出
SRC 事务开启成功,DST 事务开启失败回滚 SRC 事务
SRC 提交成功,DST 提交失败无法直接修复(没记录事务内语句),依赖后续校验

核心原则:失败就失败,等校验和修复程序找出。但需要监控双写失败频率,确保只是偶发性失败。

flowchart TD
    W[业务写请求] --> P{当前 Pattern?}
    P -- SRC_ONLY --> A[只写源库]
    P -- SRC_FIRST --> B[写源库+写目标库
目标失败吞错] P -- DST_FIRST --> C[写目标库+写源库
源失败吞错] P -- DST_ONLY --> D[只写目标库] B -. 失败记录日志 .-> M[监控双写失败频率] C -. 失败记录日志 .-> M

六、增量校验与修复

6.1 增量校验的目标

全量校验一次要几天,老数据校验过就不再校验,只校验新修改的数据

6.2 怎么知道哪些数据被修改过

思路可行性
借助双写 ConnPool 找出修改过的数据难。UPDATE xx WHERE a=1 你无法知道影响多少行
借助 utime 字段(更新时间)简单可行,但 utime 上必须有索引
借助 MySQL binlog(如 Canal)最佳方案,后续课程会讲

6.3 基于 utime 的增量校验

// IncrementalValidate 增量校验:只校验 utime > since 的数据
func (v *Validator[T]) IncrementalValidate(ctx context.Context, since int64, sleepInterval time.Duration) {
    // 步骤 1:通过 sleepInterval 控制校验节奏,避免压垮数据库
    for {
        select {
        case <-ctx.Done():
            return
        default:
        }
        var records []T
        // 步骤 2:取 utime 在 since 之后的数据
        err := v.baseDB.WithContext(ctx).
            Where("utime > ?", since).
            FindInBatches(&records, 500, func(tx *gorm.DB, batch int) error {
                // 步骤 3:逐条走通用校验逻辑
                for _, r := range records {
                    v.Validate(ctx, r.ID()-1)
                }
                return nil
            }).Error
        if err != nil {
            log.Warn("增量校验出错", err)
        }
        // 步骤 4:更新 since 为当前时间,下一轮只校验更新的
        since = time.Now().Unix()
        time.Sleep(sleepInterval)
    }
}

6.4 模式组合

utimesleepInterval行为
= 0≤ 0全量校验,结束后退出
= 0> 0全量校验,结束后继续增量校验
近期时间> 0增量校验,保持持续运行

6.5 增量校验的并发问题

目标库有两个过程在写:

  • 业务写数据(双写);
  • 修复写数据。

可能出现:修复操作直接覆盖业务最新写入

解决方案不需要办。等下次该数据被更新时,或下一轮全量/增量校验,数据会重新一致。核心原则——反复校验和修复,最终总会一致

sequenceDiagram
    participant B as 业务双写
    participant T as 目标库
    participant F as 修复写
    B->>T: 写入最新值 V2
    F->>T: 用旧值 V1 覆盖(基于旧校验快照)
    Note over T: 暂时错误
    B->>T: 下次更新 V3 / 下一轮校验发现不一致
    F->>T: 重新 UPSERT 正确值
    Note over T: 最终一致

6.6 全量校验和增量校验可以并行

  • 资源充足时可以同时开;
  • 有并发问题没关系,反复校验和修复总会一致;
  • 用 Canal 做增量时,常同时开全量校验兜底。

七、集成全部环节——Scheduler

7.1 统一调度入口

把所有组件整合到 Scheduler,提供 HTTP 接口方便操作:

// webook/internal/migration/scheduler.go
type Scheduler struct {
    validator *Validator[SomeEntity]
    fixer     *Fixer[SomeEntity]
    pool      *DoubleWritePool
}

// 全量校验
func (s *Scheduler) FullValidate() {
    go s.validator.FullValidate(context.Background())
}

// 增量校验
func (s *Scheduler) IncrementalValidate() {
    go s.validator.IncrementalValidate(context.Background(), time.Now().Unix(), 30*time.Second)
}

// 切换 pattern
func (s *Scheduler) SwitchPattern(pat Pattern) {
    s.pool.UpdatePattern(pat)
}

7.2 HTTP 管理接口

// POST /migration/full-validate
// POST /migration/incremental-validate
// POST /migration/switch-pattern?pattern=src_first
// GET  /migration/status

默认监听 8082 端口,作为管理后台端口。

7.3 演示流程

  1. 清空 webookwebook_intr 两个数据库;
  2. webook 中插入一些数据;
  3. 启动 interactive,用 wrk 持续发送请求模拟线上流量不停;
  4. 通过 Postman 触发全量校验;
  5. 全量校验通过后开启增量校验;
  6. 切换到 DST_FIRST
  7. 保持一段时间后切换到 DST_ONLY,迁移结束。

八、性能优化与数据库保护

8.1 性能优化手段

  • 按表开 goroutine:一张表一个 goroutine;
  • 按 ID 步长分片:一台机器 10 个 goroutine,分别负责尾号 0~9 的数据;
  • 多机器分布式校验:每台机器负责一部分 ID 范围。

性能瓶颈通常在数据库,不在程序本身。

8.2 数据库保护策略

  • 业务低峰期开更多 goroutine,高峰期减少甚至停掉;
  • 专门准备从库服务数据迁移(代价高,仅核心业务用);
  • 校验先读从库,不一致再读主库:从库有延迟,但大多数情况下数据是一致的,能显著降低主库压力;
  • 修复必须走主库:没办法;
  • Kafka 削峰:消费者侧限流,间接控制目标表写入速率;
  • 接入数据库限流:进一步保护主库。

8.3 主从集群下的注意事项

校验最新数据应读主库(从库有延迟)。但主库光服务写请求就负载很高,所以采用优化策略:校验先读从库,不一致再读主库。修复时只能走主库。


九、工程实践要点

  1. 顺序:先全量校验修复 → 开启双写 SRC_FIRST → 增量校验 → 切 DST_FIRST → 增量校验 → 切 DST_ONLY。
  2. 回滚:任何阶段都可以切回上一阶段,因为双写期间两边数据都在更新。
  3. 监控:双写失败频率、校验不一致数量、修复失败数量。
  4. 告警:双写失败超过阈值、校验不一致数持续增长。
  5. 人工介入:数据库查询出错只能记录日志 + 告警,由人工处理。
  6. 业务理解:异构迁移、表结构变更必须有人深度理解业务。
  7. 不要追求一次成功:反复校验和修复,最终一致即可。

十、面试要点

简历里要主动提起"主导过不停机数据迁移方案",面试中要能说清:

  • 不停机迁移的基本步骤:四阶段模型 + 8 个具体步骤。
  • 数据校验方案:泛型 Validator + Entity 实现 CompareTo + 反向校验。
  • 数据修复方案:Kafka 解耦 + UPSERT + 简化版(直接 UPSERT/DELETE,不区分类型)。
  • 如何保证数据正确性:反复校验和修复,最终一致。
  • 并发问题:业务写 vs 修复写可能互相覆盖;解决方案就是反复校验。
  • 每个阶段如何保护数据库:低峰期运行、从库校验、Kafka 削峰、限流。
  • 为什么用 Kafka:削峰 + 解耦 + 控制消费速率。
  • 主从同步下校验注意事项:先读从库,不一致再读主库。
  • 性能优化手段:按表/按 ID 分片、goroutine 并发、多机分布式。

话术模板:“在重构 XX 系统时,逼不得已要进行数据迁移,我主导设计了一个不停机迁移方案,包括双写、全量校验、增量校验、Kafka 修复、四阶段流量切换……”


自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 不停机迁移为什么叫「蜗牛悖论」?它和停机迁移的核心区别在哪?
  2. 四阶段模型里,SRC_FIRST 和 DST_FIRST 的「写」和「读」分别打向哪个库?哪个库是「基准」?
  3. 双写时 dst 写失败为什么敢「吞掉错误不返回」?靠什么机制最终把数据补回来?
  4. 修复操作为什么必须用 UPSERT(Save)而不是普通 INSERT?普通 INSERT 会在什么场景下报错?
  5. 增量校验基于 utime 有什么前提?为什么「修复覆盖业务最新写入」不需要专门解决?

动手练习(建议真做一遍)

  1. 跑通初始化:用 mysqldump 把一张小表导出成 .sql,再 source 到一个新建的 webook_xxx 库,确认数据完全一致。
  2. 手动切 Pattern:实例化一个 DoubleWritePool,分别 UpdatePattern 到 SRC_FIRST / DST_FIRST / DST_ONLY,打印每次读写实际打向的库,验证路由逻辑正确。
  3. 故意制造不一致:在 DST_FIRST 阶段,手动 UPDATE 改掉目标表某行的一个字段,触发一次校验,观察修复流程通过 Kafka 把它改回源表的值。

十一、本章小结

  1. 微服务化后必须迁移数据,核心应用采用不停机迁移。
  2. 不停机迁移核心难点:数据始终在变动(蜗牛悖论)。
  3. 四阶段模型:SRC_ONLY → SRC_FIRST → DST_FIRST → DST_ONLY,每阶段都可回滚。
  4. 双写实现首选 GORM ConnPool,能动态切换 pattern + 控制事务。
  5. 校验用泛型 Validator + Entity CompareTo,反向校验防止目标库多数据。
  6. 修复用 Kafka 解耦削峰,简化版直接 UPSERT/DELETE。
  7. 增量校验基于 utime 或 binlog(Canal),与全量校验可并行。
  8. 并发问题不解决,靠反复校验和修复最终一致。
  9. 性能瓶颈在数据库,要靠低峰期、从库、Kafka、限流多管齐下。

下章将进入服务注册与发现——微服务拆出来后,怎么让客户端找到这些动态变化的实例。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!