学习目标
学完本章,你应该能够:
- 解释为什么微服务化后必须迁移数据,说清不停机迁移的核心难点。
- 复述不停机迁移的四阶段模型:SRC_ONLY → SRC_FIRST → DST_FIRST → DST_ONLY。
- 设计数据导出导入、全量校验、增量校验、数据修复的完整方案。
- 理解基于 GORM
ConnPool的双写实现,能写出生产可用的双写代码。 - 处理迁移中的并发问题、性能优化与回滚方案,并能讲成一段面试故事。
前置知识(如果下面任意一点生疏,先回看对应章):
- 第02章 Gin + GORM:知道事务、
First/Save/批量操作怎么写。 - 第07章 Kafka:知道消息生产者/消费者、削峰与解耦(修复环节会用到)。
- 第10章 单体拆分微服务:知道「为什么要把数据库也拆开」。
- 基本的乐观锁(
version字段)、UPSERT概念。
本章你会动手做的事:
- 用
mysqldump把一张表导出来、再导入到一个新库,跑通「初始化目标表」。 - 写一个
DoubleWritePool,手动切PatternSRCFirst→PatternDSTFirst→PatternDSTOnly,观察读写打向哪个库。 - 触发一次全量校验,故意改一条目标表数据,看修复流程怎么把它改回来。
一、核心概念讲解
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 具体迁移步骤
- 创建目标表;
- 用源表数据初始化目标表;
- 执行一次校验,并用源表数据修复目标表(可选,开启双写前减少差异);
- 业务代码开启双写,读源表,先写源表,源表为准(SRC_FIRST);
- 开启增量校验和数据修复,保持一段时间;
- 切换双写顺序,读目标表,先写目标表,目标表为准(DST_FIRST);
- 继续保持增量校验和修复;
- 切换为目标表单写,读写都只操作目标表(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/Callback | Hook 要一个个结构体定义;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 模式组合
| utime | sleepInterval | 行为 |
|---|---|---|
| = 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 演示流程
- 清空
webook和webook_intr两个数据库; - 在
webook中插入一些数据; - 启动 interactive,用 wrk 持续发送请求模拟线上流量不停;
- 通过 Postman 触发全量校验;
- 全量校验通过后开启增量校验;
- 切换到
DST_FIRST; - 保持一段时间后切换到
DST_ONLY,迁移结束。
八、性能优化与数据库保护
8.1 性能优化手段
- 按表开 goroutine:一张表一个 goroutine;
- 按 ID 步长分片:一台机器 10 个 goroutine,分别负责尾号 0~9 的数据;
- 多机器分布式校验:每台机器负责一部分 ID 范围。
性能瓶颈通常在数据库,不在程序本身。
8.2 数据库保护策略
- 业务低峰期开更多 goroutine,高峰期减少甚至停掉;
- 专门准备从库服务数据迁移(代价高,仅核心业务用);
- 校验先读从库,不一致再读主库:从库有延迟,但大多数情况下数据是一致的,能显著降低主库压力;
- 修复必须走主库:没办法;
- Kafka 削峰:消费者侧限流,间接控制目标表写入速率;
- 接入数据库限流:进一步保护主库。
8.3 主从集群下的注意事项
校验最新数据应读主库(从库有延迟)。但主库光服务写请求就负载很高,所以采用优化策略:校验先读从库,不一致再读主库。修复时只能走主库。
九、工程实践要点
- 顺序:先全量校验修复 → 开启双写 SRC_FIRST → 增量校验 → 切 DST_FIRST → 增量校验 → 切 DST_ONLY。
- 回滚:任何阶段都可以切回上一阶段,因为双写期间两边数据都在更新。
- 监控:双写失败频率、校验不一致数量、修复失败数量。
- 告警:双写失败超过阈值、校验不一致数持续增长。
- 人工介入:数据库查询出错只能记录日志 + 告警,由人工处理。
- 业务理解:异构迁移、表结构变更必须有人深度理解业务。
- 不要追求一次成功:反复校验和修复,最终一致即可。
十、面试要点
简历里要主动提起"主导过不停机数据迁移方案",面试中要能说清:
- 不停机迁移的基本步骤:四阶段模型 + 8 个具体步骤。
- 数据校验方案:泛型 Validator + Entity 实现 CompareTo + 反向校验。
- 数据修复方案:Kafka 解耦 + UPSERT + 简化版(直接 UPSERT/DELETE,不区分类型)。
- 如何保证数据正确性:反复校验和修复,最终一致。
- 并发问题:业务写 vs 修复写可能互相覆盖;解决方案就是反复校验。
- 每个阶段如何保护数据库:低峰期运行、从库校验、Kafka 削峰、限流。
- 为什么用 Kafka:削峰 + 解耦 + 控制消费速率。
- 主从同步下校验注意事项:先读从库,不一致再读主库。
- 性能优化手段:按表/按 ID 分片、goroutine 并发、多机分布式。
话术模板:“在重构 XX 系统时,逼不得已要进行数据迁移,我主导设计了一个不停机迁移方案,包括双写、全量校验、增量校验、Kafka 修复、四阶段流量切换……”
自测题与动手练习
自测题(合上书能答出来,才算懂):
- 不停机迁移为什么叫「蜗牛悖论」?它和停机迁移的核心区别在哪?
- 四阶段模型里,SRC_FIRST 和 DST_FIRST 的「写」和「读」分别打向哪个库?哪个库是「基准」?
- 双写时 dst 写失败为什么敢「吞掉错误不返回」?靠什么机制最终把数据补回来?
- 修复操作为什么必须用 UPSERT(Save)而不是普通 INSERT?普通 INSERT 会在什么场景下报错?
- 增量校验基于 utime 有什么前提?为什么「修复覆盖业务最新写入」不需要专门解决?
动手练习(建议真做一遍):
- 跑通初始化:用
mysqldump把一张小表导出成.sql,再source到一个新建的webook_xxx库,确认数据完全一致。 - 手动切 Pattern:实例化一个
DoubleWritePool,分别UpdatePattern到 SRC_FIRST / DST_FIRST / DST_ONLY,打印每次读写实际打向的库,验证路由逻辑正确。 - 故意制造不一致:在 DST_FIRST 阶段,手动
UPDATE改掉目标表某行的一个字段,触发一次校验,观察修复流程通过 Kafka 把它改回源表的值。
十一、本章小结
- 微服务化后必须迁移数据,核心应用采用不停机迁移。
- 不停机迁移核心难点:数据始终在变动(蜗牛悖论)。
- 四阶段模型:SRC_ONLY → SRC_FIRST → DST_FIRST → DST_ONLY,每阶段都可回滚。
- 双写实现首选 GORM ConnPool,能动态切换 pattern + 控制事务。
- 校验用泛型 Validator + Entity CompareTo,反向校验防止目标库多数据。
- 修复用 Kafka 解耦削峰,简化版直接 UPSERT/DELETE。
- 增量校验基于 utime 或 binlog(Canal),与全量校验可并行。
- 并发问题不解决,靠反复校验和修复最终一致。
- 性能瓶颈在数据库,要靠低峰期、从库、Kafka、限流多管齐下。
下章将进入服务注册与发现——微服务拆出来后,怎么让客户端找到这些动态变化的实例。