作者:Oliver Nguyen

作为一家初创公司,我们最初采用的架构十分简单:所有数据都存储在一个由几个服务共享的中心化 Postgres 实例中。它一直运行得很好,让我们能够快速推进、交付功能、进入市场、获取客户并实现指数级增长。
时间快进几年,如今我们有着数千个客户和数十亿条记录,是时候把最核心的几张表从 Postgres 迁移到更好的方案了。我们选择了 DynamoDB,一款 AWS 提供的键值数据存储,具有高可用性、低延迟性能,查询延迟降到了个位数毫秒级别。
和任何靠谱的迁移方案一样,我们仔细分析了使用场景,定义了数据模型和索引,从双写开始,逐步把读查询切换到 DynamoDB,Postgres 作为回退,把所有记录导入 DynamoDB,最后停止查询 Postgres。
这篇文章聚焦于导入 room_events 表记录这一步骤:扫描该表的所有记录并写入 DynamoDB*。目标是确保迁移完整,不遗漏任何事件。迁移脚本必须够快,能够并行运行,并且能够从上次的状态中断点续跑。它还需要具备韧性,能够处理网络错误或中断,并无缝地恢复运行。*
来看看我们是怎么做的! 🥸
以下是 Postgres 中 room_events 表的简化数据模型:
room_events
说明:
我们先写一个简单版本的迁移脚本:在单个循环中运行,扫描所有记录,并写入 DynamoDB。同时加入一些琢磨透了的琐碎细节,免得之后再操心:
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(/* ... */ ) { /* ... */ }
这个朴素的脚本建立在一些假设之上:
然而,在真实世界里,我们无法依赖这些假设:
为了应对这些挑战,我们需要改进脚本,实现一套更稳健的分布式方案:
数据拆分与分区:
用多个 Worker 实现并行处理:
在 Redis 中管理状态:
用 Keeper 进行集中式进度追踪:
超时处理与重试机制:
数据校验:
利用写入的幂等性:
另一个重要的点是,每条记录的写入步骤是幂等的。这意味着把同一条记录多次写入 DynamoDB,只会简单地覆盖之前的版本,确保不会产生重复数据。
我们可以利用这一点:
有了这些概念,让我们开始思考架构设计。
既然我们已经把记录拆分成了 Timeslot,每个时间片对应一个小时内的所有 room events。我们需要在 Redis 中追踪这些时间片的进度:
所以每个时间片需要 2 个 key。乘以 3 年的数据量,会有 52,560 个 key(3 * 365 * 24 * 2)。数量相当可观!
我们可以通过清理连续一段 FINISHED 的时间片、并用单个 mgre:last_slot key 替代它们,来优化 Redis 的 key 数量。这样一来,我们只需要维护少量的 key。当然,这基于一个假设:各个时间片的 room events 数量相近,因此每个 Worker 处理完的耗时也相近。
脚本运行在多个 Pod 中。每个 Pod 包含一个 Manager、一个 Keeper 和多个 Worker,各自运行在自己的 goroutine 中。
Manager 管理一个 Pod 中的所有 Worker:
Worker 按每个 Timeslot 迁移 room events:
Keeper 总览整个迁移进度:
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 中,并会持续自我重试:
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:
用 context.WithCancel() 停止所有任务:
把进度保存为状态,并定期刷新锁:
Manager 定期检查 Redis 中的 mgre:last_slot:
记录并汇报进度:
监控资源并限流:
每个 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)}
}
正如前面讨论过的,我们有 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...)
}
}
这里有许多状态:全局状态、全局进度、时间片状态、时间片进度、锁等等。对于每一种状态,我们都需要刷新、更新或删除。每个操作也都需要能够重试:
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"))
// ...
}
好多了!
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。