Connectly
技术2024-10-22

用 Redis 实现分布式状态机,迁移数十亿条记录

作者:Oliver Nguyen

用 Redis 实现分布式状态机,迁移数十亿条记录

作为一家初创公司,我们最初采用的架构十分简单:所有数据都存储在一个由几个服务共享的中心化 Postgres 实例中。它一直运行得很好,让我们能够快速推进、交付功能、进入市场、获取客户并实现指数级增长。

时间快进几年,如今我们有着数千个客户和数十亿条记录,是时候把最核心的几张表从 Postgres 迁移到更好的方案了。我们选择了 DynamoDB,一款 AWS 提供的键值数据存储,具有高可用性、低延迟性能,查询延迟降到了个位数毫秒级别。

和任何靠谱的迁移方案一样,我们仔细分析了使用场景,定义了数据模型和索引,从双写开始,逐步把读查询切换到 DynamoDB,Postgres 作为回退,把所有记录导入 DynamoDB,最后停止查询 Postgres。

这篇文章聚焦于导入 room_events 表记录这一步骤:扫描该表的所有记录并写入 DynamoDB*。目标是确保迁移完整,不遗漏任何事件。迁移脚本必须够快,能够并行运行,并且能够从上次的状态中断点续跑。它还需要具备韧性,能够处理网络错误或中断,并无缝地恢复运行。*

来看看我们是怎么做的! 🥸

数据模型

以下是 Postgres 中 room_events 表的简化数据模型:

room_events

  • id:ULID,主键
  • room_id:UUID,引用 rooms.id
  • data:JSON

说明:

  • id 列是 ULID,带有时间成分。
  • 该表拥有数十亿条记录。

一个朴素的方案

我们先写一个简单版本的迁移脚本:在单个循环中运行,扫描所有记录,并写入 DynamoDB。同时加入一些琢磨透了的琐碎细节,免得之后再操心:

  • 基于 id 列使用游标分页进行扫描。
  • 每次写入使用最多 25 条记录的批次(BatchWriteItem 允许的最大条目数)。
  • 从双写开始的日期往最早的日期反向扫描。这样我们能优先获得最新的记录。
import "connectly.ai/go/pkgs/ulid"

var startDate = parseTime("2020-01-01T00:00:00Z")
var endDate   = parseTime("2024-10-20T00:00:00Z")
var pgBatchSize  = 1000
var dynBatchSize = 25

type Migrator struct {
    *logger
    pg  PostgresClient
    dyn DynamoClient
}
type QueryRoomEventsRequest struct {
    BeforeID ulid.ULID
    Limit    int
}
type QueryRoomEventsResponse struct {
    Items  []RoomEvents
    lastID ulid.ULID
}

func (m *Migrator) Migrate(ctx context.Context) error {
    lastID, count := ulid.FromTime(endDate), 0
    for {
        req := QueryRoomEventsRequest{
            BeforeID: lastID,
            Limit:    pgBatchSize
        }
        res, err := m.queryRoomEvents(ctx, req)
        m.Logger(ctx).Must(err, "failed to query room events")
        count += res.Items
        if len(res.Items) == 0 {
            m.Logger(ctx).Infof("MIGRATION DONE: count=%v", count)
            return nil
        }
        for i := 0; i < len(res.Items); i += dynBatchSize {
            j := min(i + dynBatchSize, len(res.Items))
            batch := res.Items[i, j]
            _, err = m.batchWriteDynamo(ctx, batch)
            m.Logger(ctx).Must(err, "failed to write to dynamodb")
        }
    }
}

func (m *Migrator) queryRoomEvents( /* ... */ ) { /* ... */ }
func (m *Migrator) batchWriteDynamo(/* ... */ ) { /* ... */ }

真实世界的方案

这个朴素的脚本建立在一些假设之上:

  • 对 queryRoomEvents() 和 batchWriteDynamo() 的调用总是成功。
  • 记录数量足够小,可以在单台机器的单次运行中完成。

然而,在真实世界里,我们无法依赖这些假设:

  • 请求可能随机超时或失败。
  • 数据库可能触及吞吐量或连接数上限。
  • Pod 可能因为故障或新部署而重启。
  • 海量的记录会让单线程方案变得不切实际,需要耗费极长时间才能完成。

为了应对这些挑战,我们需要改进脚本,实现一套更稳健的分布式方案:

数据拆分与分区:

  • 把 room events 按 1 小时的 Timeslot(时间片)分组拆分。每个时间片应代表一个可独立处理的、可控的工作单元。这样可以实现更好的负载均衡和更强的并行性。

用多个 Worker 实现并行处理:

  • 在多个 Pod 上运行多个 Worker,每个 Pod 管理多个 goroutine。
  • 每个 Worker 领取一个 Timeslot,完成迁移后再领取下一个可用的时间片,确保所有 Worker 都能并发工作。

在 Redis 中管理状态:

  • Worker 在处理某个 Timeslot 之前,需要先在 Redis 中获取该时间片的锁,确保同一时刻只有一个 Worker 处理该时间片。
  • 处理完一个 Timeslot 后,Worker 会把该时间片的状态更新为 FINISHED。
  • 如果某个 Worker 失败,锁会被释放或过期,让另一个 Worker 接管并继续处理。

用 Keeper 进行集中式进度追踪:

  • 引入一个 Keeper 来监控所有 Worker 的进度,维护整个迁移的全局概览,并记录全局进度。
  • 它会追踪每个 Timeslot 的进度和状态,确保没有任何一个时间片被遗漏。

超时处理与重试机制:

  • 为 PostgreSQL 读取和 DynamoDB 写入都实现带指数回退的重试机制,让脚本能够抵御瞬时网络故障或超时。
  • 再增加一层重试和错误处理,防止任何 panic 或意外故障,确保脚本能够优雅地恢复,不丢数据、不产生损坏的状态。

数据校验:

  • 迁移完成后,运行一个校验脚本,确保所有 room events 都已成功迁移,没有遗漏的记录。
  • 校验过程应比对 PostgreSQL 和 DynamoDB 之间的记录数量以及样本数据。

利用写入的幂等性:

另一个重要的点是,每条记录的写入步骤是幂等的。这意味着把同一条记录多次写入 DynamoDB,只会简单地覆盖之前的版本,确保不会产生重复数据。

我们可以利用这一点:

  • 如果在脚本中发现了 bug,需要运行新版本,我们可以重置所有状态,从头开始运行,而不必删除 DynamoDB 中现有的数据。
  • 如果脚本因为某种原因停止后又恢复运行,它可能会覆盖 DynamoDB 中的一些记录。这是可以接受的,偶尔的重写是无害的,能避免重新做大量已经完成的工作。

架构

有了这些概念,让我们开始思考架构设计。

在 Redis 中存储状态

既然我们已经把记录拆分成了 Timeslot,每个时间片对应一个小时内的所有 room events。我们需要在 Redis 中追踪这些时间片的进度:

  • 所有 key 都加上 mgre: 前缀 → 便于扫描和批量删除。
  • mgre:status → 全局迁移状态:NOT_STARTED、TO_START、IN_PROGRESS、TO_STOP、STOPPED、FINISHED。
  • mgre:{TIMESLOT}:worker → 当前正在处理该时间片的 Worker,一个 TTL 较短的独占锁(例如 15-30 秒)。当某个 Worker 未能释放锁时,锁会自动过期,让另一个 Worker 稍后接管。
  • mgre:{TIMESLOT}:states → 一个存储该时间片状态的 JSON,TTL 较长(例如 1 周),包含以下属性:
  • .status → 空、IN_PROGRESS 或 FINISHED。注意这里没有 ERROR 状态。所有 room events 和所有时间片都必须成功迁移! 😎
  • .last_id → 最后一条已迁移的 room events id,用于恢复进度。

所以每个时间片需要 2 个 key。乘以 3 年的数据量,会有 52,560 个 key(3 * 365 * 24 * 2)。数量相当可观!

我们可以通过清理连续一段 FINISHED 的时间片、并用单个 mgre:last_slot key 替代它们,来优化 Redis 的 key 数量。这样一来,我们只需要维护少量的 key。当然,这基于一个假设:各个时间片的 room events 数量相近,因此每个 Worker 处理完的耗时也相近。

Manager、Keeper 与 Worker

脚本运行在多个 Pod 中。每个 Pod 包含一个 Manager、一个 Keeper 和多个 Worker,各自运行在自己的 goroutine 中。

Manager 管理一个 Pod 中的所有 Worker:

  • 每个 Pod 只有一个 Manager。
  • 它遍历所有的 Timeslot。
  • 对每个时间片,它通过检查 mgre:{TIMESLOT}:states 和 mgre:{TIMESLOT}:worker 来验证对应的状态和锁是否存在于 Redis 中。
  • 然后把该时间片派发给 Worker,方式是把它发送到一个 Go channel slotCh 中。

Worker 按每个 Timeslot 迁移 room events:

  • 有多个 Worker,每个运行在独立的 goroutine 中。
  • Worker 从 channel slotCh 中一次接收一个 Timeslot。
  • 它再次在 mgre:{TIMESLOT}:states 中验证状态。
  • 通过在 mgre:{TIMESLOT}:worker 上执行 SET NX 命令来获取锁。
  • 如果一切成功,就开始处理该时间片。
  • 它会不时地通过更新 mgre:{TIMESLOT}:states 来保存进度。
  • 并能够从已保存的进度中恢复。
  • 如果失败,锁最终会过期,让另一个 Worker 接管。

Keeper 总览整个迁移进度:

  • 所有 Pod 中应该只有一个处于活跃状态的 Keeper。
  • 我们应该有一个 TTL 较短的锁 mgre:keeper,这样一旦某个 Pod 上的 Keeper 失败,另一个 Keeper 就能接管。
  • 一个失败的 Keeper 会回到待命状态,并继续监视该锁,看自己是否能够重新变为活跃状态。
  • 每隔几秒,处于活跃状态的 Keeper 会从 mgre:last_slot 开始扫描,找出连续的 FINISHED Timeslot。然后它会把 mgre:last_slot 标记为最新的连续 FINISHED 时间片,并清理它们。
  • 当所有剩余的时间片都变为 FINISHED 时,它会把全局 mgre:status 标记为 FINISHED。所有 Manager 都会关注这个状态,一旦看到就会停止运行。
func (m *Migrator) runManager(ctx context.Context) {
    // ... recover, logs, metrics ...
    // ... loop and call runManagerStep (more on this later) ...
}

func (m *Migrator) runManagerStep(ctx context.Context) Status {
    // ... recover, logs, metrics ...
    
    for {
        mustRefresh(ctx, m, m.globalStatus, m.lastSlot)
        if m.globalStatus.Load().Is(SKIPPED, FINISHED) { return /*...*/ }
        
        for slot := m.lastSlot.Load(); slot.Before(endSlot); slot.Next() {
            slotStates := m.getSlotStates(slot)
            mustRefresh(ctx, m, slotStates)
            if slotStates.Status.Load().Is(SKIPPED, FINISHED) { continue }
            
            slotWorkerID := mustAcquireLock(ctx, m, slot.Worker)
            if !slotWorkerID.Load().IsPod(m.PodID) { continue }

            select {
                case <- ctx.Done():
                    return IN_PROGRESS  // 👈 stop when the pod stops
                case m.slotCh <- slot: 
                    continue            // 👈 send to slot channel
            }
        }
    }
    // 👇 still IN_PROGRESS, only Keeper can verify all slots are FINISHED
    return IN_PROGRESS
}

func (m *Migrator) runWorker(ctx context.Context) {
    // ... recover, logs, metrics ...
    for {
        select {
            case <- ctx.Done():
                return
            case slot := <- m.slotCh
                s.runTimeslot(ctx, slot)
        }
    }
}

func (m *Migrator) runTimeslot(ctx context.Context, slot Timeslot) {
    // ... recover, logs, metrics ...
    t := time.NewTimer(0)
    for {
        select {
            case <- ctx.Done():
                return
            case <- t.C
                status := s.runTimeslotStep(ctx, slot)
                if status.Is(SKIPPED, FINISHED) { return }
                t.Reset(3 * time.Second)
        }
    }
}

func (m *Migrator) runTimeslotStep(ctx context.Context, slot Timeslot) Status {
    // ... recover, logs, metrics ...
    st := m.getSlotStates(slot)
    for {
        // 👇 refresh the states
        mustRefresh(ctx, m, st.States)
        // 👇 verify the status, return if already FINISHED
        if st.States.Status.Load().Is(FINISHED) { return FINISHED }
        // 👇 another worker is working on the slot, SKIPPED
        if !mustAcquireLock(ctx, m, st.Worker, workerID) { return SKIPPED }
        lastID := st.States.LastID
        for lastID.Before(startID) {
            req := m.QueryRoomEventsRequest{ BeforeID: lastID; Limit: ... }
            res := mustRetry(ctx, m.strategy, msgf("query room events"), 
                func() (QueryRoomEventsResponse, error) { 
                    return m.queryRoomEvents(ctx, req)
                }
            // 👉 ... save to DynamoDB
            // 👉 ... save progress to Redis
            lastID = res.LastID
        } 
    }
}

func (m *Migrator) runKeeper(ctx context.Context) {
    // ... recover, logs, metrics ...
    // ... loop and call runKeeperStep (more on this later) ...
}

func (m *Migrator) runKeeperStep(ctx context.Context) Status {
    // ... recover, logs, metrics ...
    for {
        // 👇 acquire lock and refresh states
        //    only a single active Keeper across all pods
        if !mustAcquireLock(ctx, m, m.keeper) { return SKIPPED }
        mustRefresh(ctx, m, m.globalStatus, m.lastSlot)
        if m.globalStatus.Load().Is(FINISHED) { return FINISHED }
        // 👇 find the last consecutive FINISHED slot
        newLastSlot := states.LastSlot
        for slot := states.LastSlot; slot.Before(endSlot); slot.Next() {
            slotStates := m.getSlotStates(slot)
            err := tryRefresh(ctx, m, slotStates)
            if err != nil { break }
            if !slotStates.Status.Load().Is(FINISHED) { break }
            
            newLastSlot = slot
            tryClean(ctx, m, slotStates) // 👈 clean FINISHED slot
        }
        // 👇 save the mgre:last_slot state
        if newLastSlot != states.LastSlot {
            mustUpdate(ctx, m, m.lastSlot, newLastSlot)
        }
    }
}

重试机制

计划看起来不错,对吧?不,还没完。如果上面任何一个步骤失败了会怎样?

每个 Manager、 Worker、 Keeper 都运行在一个 goroutine 中,并会持续自我重试:

  • 始终使用 recover(),因为 goroutine 中任何未被 recover 的 panic 都会中止整个进程。
  • 有 2 个主循环,用于在出现问题时能够重新启动:
  • 外层循环负责重启内层循环。它只包含简单的语句,确保自身永远不会 panic。
  • 内层循环负责处理业务逻辑:加载状态、查询数据等等。
  • 对于 Worker,还有一层最外层的循环,用于接收时间片,并将其传递给外层循环再传给内层循环去执行。
  • 还有一个重试函数,用于执行并对每次查询多重试几次。
func (m *Migrator) initAndRun(ctx context.Context) {
    go m.runManager(ctx)
    go m.runKeeper(ctx)
    for i := 0; i < numWorkers; i++ {
        go m.runWorker(ctx)
    }
}

func (m *Migrator) runWorker(ctx context.Context) {
    defer func() {
        r := recover()
        if r != nil { log(ctx).Errorf("panic in the outermost layer, will stop") }
    }
    // 👇 the outermost loop to receive the next slot
    //    it only contains simple statements to ensure that it never panics
    for {
        select {
            case <- ctx.Done():
                return
            case slot := <- m.slotCh  // 👈 receive slots from channel
                s.runTimeslot(ctx, slot) // and execute them one by one
        }
    }
}

func (m *Migrator) runTimeslot(ctx context.Context) {
    defer func() {
        r := recover()
        if r != nil { log(ctx).Errorf("panic in the outer layer, will stop") }
    }
    // 👇 the outer loop to retry the migration logic
    //    it only contains simple statements to ensure that it never panics
    t := time.NewTimer(0)
    for {
        select {
            case <-ctx.Done(): 
                return                // 👈 stop when the pod stops
            case <-t.C:
                status := m.runTimeslotStep()
                if status.Is(FINISHED, SKIPPED) { 
                    return            // 👈 stop when FINISHED or SKIPPED
                }
                t.Reset(3 * time.Second) // 👈 retry after a few sec
        }
    }
}

func (m *Migrator) runTimeslotStep(ctx context.Context) (Status) {
    defer func() {
        r := recover()
        if r != nil { logger(ctx, "panic in the inner layer, will retry") }
    }
    // 👇 the inner loop to execute the migration logic
    for {
        // ... load states, progress, acquire lock...
        // 👇 query database
        res, err := retry(ctx, m.strategy, msgf("query room events"), 
            func() (QueryRoomEventsResponse, error) {
                return m.queryRoomEvents(/* ... */)
            })
        // 👇 even if there is panic, the runManagerStep will recover, stop
        //    and the outer loop (runManager) will continue retry after few sec
        must(err)
        
        // ... save states, progress, refresh lock
    }
    // ... save status as FINISHED
    return FINISHED // 👈 tell the outer loop to stop
}

当出现任何错误时,例如网络超时:

  • 重试函数会重试几次。
  • 如果失败,内层循环会停止,把控制权交还给外层循环。
  • 外层循环会等待几秒钟,然后重新启动内层循环。
  • 内层循环随后加载之前的状态和进度,获取锁,并继续执行。

这确保了代码会一直运行,直到所有记录都迁移完成,或者 Pod 重启。在后一种情况下,它会在下次运行时恢复迁移进度。

其他说明

用于控制迁移的 API:

  • 我们应该暴露一个 API 来控制迁移:启动和停止所有 Pod 上的 Manager、Keeper、Worker。
  • 收到 START 命令时,该 API 会把 mgre:status 设为 TO_START 或 IN_PROGRESS。
  • 当所有 Pod 上的 Manager 看到这个状态时,它们就会启动或恢复迁移。

用 context.WithCancel() 停止所有任务:

  • 使用单个 context.Context,并把它传递给同一个 Pod 中的所有 manager、keeper 和 worker。
  • 当 Manager 看到状态变为 TO_STOP、STOPPED 时,它们会取消该 context,让 keeper 和所有 worker 优雅地停止。

把进度保存为状态,并定期刷新锁:

  • 当每个 worker 或 keeper 在处理各自的任务时,应该把进度保存到 Redis 中,并刷新锁,以便之后能够恢复,并防止其他 worker 抢占该时间片。

Manager 定期检查 Redis 中的 mgre:last_slot:

  • 即便某个时间片已经被某个 worker 领取,仍有可能出现该 worker 停止而这个时间片还没有变成 FINISHED 的情况。
  • 每个 Manager 在遍历所有时间片的过程中,过一段时间后,需要重新开始这个循环,检查 Redis 中 mgre:last_slot 的状态,然后从那里继续。这是为了确保我们不会漏下任何未完成的时间片。

记录并汇报进度:

  • Keeper 负责总览进度、清理已完成的记录,并更新 mgre:last_slot。
  • 在任务进行期间,Keeper 应该定期记录并汇报进度,以便实时了解迁移状态,并对整体进度进行实时追踪。

监控资源并限流:

  • 迁移过程可能会占满数据库的所有资源。我们可能需要监控并添加限流机制,并在必要时微调参数,确保一切尽可能顺利地进行。

实现

Timeslot

每个 Timeslot 代表一小时内的所有 room events。我们可以把它实现为一个 time.Time,并以字符串格式(如 20241020.02)存储在 Redis 中。

type Timeslot struct { time.Time }

const slotDuration = time.Hour
func newTimeSlot(t time.Time) Timeslot {
    if t.IsZero() { return Timeslot{t} }
    // 👉 each slot is an hour
    t = t.In(time.UTC).Truncate(kSlotDuration)
    return Timeslot{t}
}
func (t Timeslot) String() string {  
    if t.Time.IsZero() { return "" }  
    return t.Time.Format("20060102.15")  
}
func (t Timeslot) Range() (start, end ulid.ULID) {
    return ulid.FromTime(t), ulid.FromTime(t.Add(1))
}
func (t Timeslot) Add(i int) Timeslot {
    return Timeslot{t.Time.Add(slotDuration)}
}
func (t Timeslot) Next() Timeslot {
    return t.Add(-1) // 👈 we are going from latest to earliest
}
func (t Timeslot) Sub(x Timeslot) int {
    return Timeslot{t.Time.Sub(kSlotDuration)}
}

Retry() 查询

正如前面讨论过的,我们有 2 个循环来处理 panic、重试和状态恢复。所以对于这个 retry() 函数,我们只需要把查询重试几次,就能容忍一些网络故障:

// 👉 this will retry 3 times
strategy := NewSimpleStrategy(
    100*time.Millisecond, 200*time.Millisecond, 500*time.Millisecond)

// 👉 call the QueryRoomEvents with retry-ability
func retry(ctx, strategy, msgf("query database"),
    func() (QueryRoomEventsRequest, error) {
        return m.QueryRoomEvents(ctx, req)
    })
// 👉 if all retries failed, stop, and give control back to outer loop
//    to try again after a few sec
func mustRetry( /* ... */ ) { /* ... */ }

我们可以实现一个简单的重试逻辑:

func retry[T any](
    ctx context.Context, strategy RetryStrategy,
    msg fmt.Stringer, fn func() (T, error),
) (T, error) {
    for count := 0; ; count++ {
        x, err := fn()
        if err == nil { return x, err }
        if next := strategy.Next(); next > 0 {
            logger(ctx).Warnf("failed to %v (attempt %v)", msg, count)
            time.Sleep(next)
        } else {
            logger(ctx).Errorf("failed to %v (attempt %v)", msg, count)
            return x, err
        }
    }
}

重试策略的一个示例实现,每次重试会在预设的间隔后发生:

type RetryStrategy func() time.Duration

func (f RetryStrategy) Next() time.Duration { return f() }
func NewSimpleStrategy(at []time.Duration) RetryStrategy {
    i := -1
    return func() time.Duration {
        i++
        if i < len(at) { return i }
        return -1; // no more retry
    }
}

以及 msgf() 函数的实现,用于快速返回一个 fmt.Stringer:

type StringFunc func() string

func (f StringFunc) String() { return f() }
func msgf(msg string, args ...any) StringFunc {
   return func() string {
       return fmt.Sprintf(msg, args...)
   }
}

States(状态)

这里有许多状态:全局状态、全局进度、时间片状态、时间片进度、锁等等。对于每一种状态,我们都需要刷新、更新或删除。每个操作也都需要能够重试:

globalStatus := mustRetry(ctx, m.strategy, msgf("load status")
	func() (Status, error) {
	    str, err := m.redisClient.GetString(ctx, "mgre:status")
	    if err != nil { return 0, err }
	    if str == "" { return NOT_STARTED }
	    return parseStatus(str)
	})
slotWorker := mustRetry(ctx, m.strategy, msgf("load last slot")
	func() (string, error) {
	    return m.redisClient.GetString(ctx, "mgre:", kSlotWorker(slot))
	})	
mustRetry(ctx, m.strategy, msg("save status")
	func() (int, error) {
	    str := encodeStatus(newStatus)
	    err := m.redisClient.SetStringTTL(ctx, "mgre:status", str)
	    return 0, err
	})

func kSlotWorker(slot fmt.Stringer) string {
    return fmt.Sprintf("mgre:%v:worker", slot)
}

这样的代码很快会变得过于冗长。我们可以把 key(包括编码逻辑和其他配置)封装进一个 State 结构体中:

type State[T any] struct {
    v      atomic.Value
    key    string
    ttl    time.Duration
    retry  RetryStrategy
    parse  func(string) (T, error)
    encode func(T) (string, error)
}

func NewState[T any](
    key string, ttl time.Duration, strategy RetryStrategy,
    parse  func(string) (T, error), 
    encode func(T) (string, error) 
) *State[T] {
    var zero T
    st := &State[T]{ /* ... */ }
    st.v.Store(zero)
    return st
}
func (s *State[T]) Load() T {
    return s.v.Load().(T)
}
func (s *State[T]) Refresh(ctx context.Context, redis RedisClient) error {
    str, err:= retry(ctx, msgf("get key %q", s.key),
        func() (string, error) {
            return redisClient.Get(ctx, key)
        })
    if err != nil { return err }
    v, err := encode(str)
    if err != nil { return err }
    s.v.Store(v)
    return nil
}
func (s *State[T]) Save(ctx context.Context, redis RedisClient, v T) error {
    str, err := s.encode(v)
    if err != nil { return err }
    return retry(ctx, msgf("set key %q", s.key), 
        func() (int, error) {
            _, err := redisClient.Set(ctx, key, str, s.ttl)
            return 0, err
        })
}
func (s *State[T]) AcquireLock(ctx context.Context, redis RedisClient, v T) error {
    // 👉 similar to Save(), use SetNX instead ...
}

再实现一些辅助函数以便快速访问它们:

type DepsI  interface { _redis() RedisClient }
type StateI interface { _key() string; _ttl() time.Duration; /* ... */ }

func mustRefresh(ctx context.Context, deps DepsI, states StateI, msg fmt.Stringer) {
    /* ... */
}
func mustSet[T any](ctx context.Context, deps DepsI, state State[T], v T) {
    /* ... */
}

最后,我们就可以简化用法了:

func (m *Migrator) exampleInit(slot Slot) {
    m.globalStatus = NewState(kStatus, longTTL, m.strategy, parseStatus, encodeStatus)
    m.slotWorker = NewState(kSlotWorker(slot), shortTTL, m.strategy, parseStr, encodeStr)
    // ...
}

func (m *Migrator) exampleRefreshStates() {
    mustRefresh(ctx, m, []StateI{m.globalStatus, m.lastSlot}, msgf("load states"))
    mustSet(ctx, m, m.globalStatus, FINISHED, msgf("save status"))
    // ...
}

好多了!

SlotStates

Timeslot 结构体只定义了时间片本身。我们还需要另一个结构体来封装它的状态:

type SlotStates struct {
    Timeslot             // 👉 embedded Timeslot to quickly access methods
    status State[Status] // 👉 mgre:TIMESLOT:status
    worker State[string] // 👉 mgre:TIMESLOT:worker
}

func newSlotStates(slot Timeslot) *SlotStates {
    return &SlotStates{
        Timeslot: slot,
        status: NewState(kStatus, longTTL,  /* ... */),
        worker: NewState(kWorker, shortTTL, /* ... */),
    }
}

在一个 Pod 中,每个 Timeslot 应该只有一份 SlotStates,由 manager、keeper 和 worker 共享。所以最好有一个集中的地方来初始化并存储它们:

type Migrator {
    // ...
    slots map[string]*SlotStates
    mu    sync.RWMutex
}

func (m *Migrator) getSlotStates(slot Timeslot) *SlotStates {
    if st := m._getSlotStates(); st != nil { return st }
    m.mu.Lock()
    defer m.mu.Unlock()
    if st := m.slots[slot.String()]; st != nil { return st }
    st := newSlotStates(slot)
    s.slots[slot.String()] = st
    return st
}
func (m *Migrator) _getSlotStates() *SlotStates {
    m.mu.RLock()
    defer m.mu.RUnlock()
    return m.slots[slot.String()]
}

结语

呼!这内容可真不少!很高兴你还在坚持读完!

在真实环境中迁移海量数据,远比一个简单脚本能处理的复杂得多。通过精心设计和实现迁移代码——包括分区、并行处理、状态管理、故障容错、幂等性以及集中式进度追踪——我们才能实现一个可靠的迁移流程,保证数据完整性,并把停机时间降到最低。

而且也能让人晚上睡得安稳!😋

作者

我是 Oliver Nguyen —— Connectly.ai 的一名软件工程师。我喜欢不断学习,每天都希望看到更好的自己。偶尔会分拆出一些新的开源项目。在旅程中分享知识和想法。

本文同时发布于 olivernguyen.io。