Por Oliver Nguyen

Como uma startup, começamos com uma arquitetura simples: todos os dados eram armazenados em uma única instância centralizada do Postgres, compartilhada por alguns serviços. Isso funcionou bem e nos permitiu avançar rápido, entregar features, entrar no mercado, conquistar clientes e crescer exponencialmente.
Alguns anos depois, hoje temos milhares de clientes e bilhões de registros; é hora de migrar as principais tabelas do Postgres para uma solução melhor. Escolhemos o DynamoDB, um datastore de chave-valor da AWS, com alta disponibilidade, baixa latência, e vimos as consultas caírem para milissegundos de um único dígito.
Como em qualquer bom plano de migração, analisamos cuidadosamente os usos, definimos schema e índices, começamos com double-write, trocamos as queries de leitura para o DynamoDB com fallback para o Postgres, importamos todos os registros para o DynamoDB e, por fim, paramos de consultar o Postgres.
Este artigo foca na etapa de importar registros da tabela room_events: escanear todos os registros dessa tabela e gravá-los no DynamoDB*. O objetivo é garantir uma migração completa, sem perder nenhum evento. O script de migração precisa ser rápido, capaz de rodar em paralelo, e capaz de parar e retomar a partir do seu último estado. Também precisa ser resiliente, lidando com erros de rede ou interrupções e retomando de forma transparente.*
Vamos ver como fazer isso! 🥸
Aqui está o schema simplificado da tabela room_events no Postgres:
room_events
Observações:
Vamos escrever uma versão simples do script de migração: rodar em um único loop, escanear todos os registros e gravar no DynamoDB. Também vamos incluir algumas coisas triviais para não precisarmos nos preocupar com elas depois:
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(/* ... */ ) { /* ... */ }
Este script ingênuo vai operar sob algumas premissas:
Porém, no mundo real, não podemos confiar nessas premissas:
Para lidar com esses desafios, precisamos melhorar nosso script e implementar uma solução distribuída mais robusta:
Divisão e particionamento dos dados:
Processamento paralelo com múltiplos workers:
Gerenciamento de estados no Redis:
Keeper para acompanhamento centralizado do progresso:
Tratamento de timeout e mecanismo de retry:
Verificação de dados:
Aproveitando a idempotência nas gravações:
Outro ponto importante é que a etapa de gravação de cada registro é idempotente. Isso significa que gravar o mesmo registro no DynamoDB várias vezes simplesmente sobrescreve a versão anterior, garantindo que não haja duplicatas.
Podemos usar isso a nosso favor:
Com esses conceitos em mente, vamos começar a pensar na arquitetura.
Dado que já dividimos os registros em Timeslots, cada um consistindo em todos os room events de uma hora, precisamos acompanhar o progresso desses slots no Redis:
Então precisamos de 2 chaves para cada slot. Multiplicando por 3 anos de dados, teremos 52.560 chaves (3 * 365 * 24 * 2). Isso é bastante!
Podemos otimizar as chaves do Redis limpando o grupo de slots consecutivos FINISHED e substituindo-os por uma única chave mgre:last_slot. Assim, só precisamos acompanhar uma pequena quantidade de chaves. Claro, isso se baseia na premissa de que os slots têm um número similar de room events, então cada worker leva um tempo similar para terminar.
O script roda em múltiplos pods. Cada pod consiste em um Manager, um Keeper e vários Workers. Cada um roda na sua própria goroutine.
Um Manager gerencia todos os workers em um pod:
Um Worker migra os room events de cada Timeslot:
Um Keeper monitora todo o progresso da migração:
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)
}
}
}
O plano parece bom, certo? Não, ainda não. O que acontece se algum dos passos acima falhar?
Cada Manager, Worker e Keeper roda em uma goroutine e sempre se autorretenta:
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
}
Quando qualquer erro acontece, por exemplo, um timeout de rede:
Isso garante que o código sempre continue rodando até que todos os registros sejam migrados, ou até que o pod reinicie. Nesse último caso, ele vai retomar o progresso da migração na próxima vez.
API para controlar a migração:
Parar tudo usando context.WithCancel():
Salvar o progresso como estados e renovar o lock periodicamente:
O Manager verifica periodicamente o mgre:last_slot no Redis:
Registrar e reportar o progresso:
Monitorar os recursos e limitar a taxa:
Cada Timeslot representa todos os room events em uma hora. Podemos implementá-lo como um time.Time e armazená-lo no Redis como uma string no formato 20241020.02.
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)}
}
Como discutido antes, temos 2 loops para lidar com pânico, retry e retomada de estados. Então, para esta função retry(), só precisamos tentar a consulta algumas vezes, para conseguir tolerar algumas falhas de rede:
// 👉 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( /* ... */ ) { /* ... */ }
Podemos implementar uma lógica de retry simples:
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
}
}
}
Um exemplo de implementação de estratégia de retry, em que cada tentativa acontece após uma duração pré-definida:
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
}
}
E a implementação da função msgf(), que pode ser usada para rapidamente retornar um 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...)
}
}
Existem muitos estados: status global, progresso global, status do slot, progresso do slot, locks, etc. Para cada estado, precisamos atualizar, salvar ou deletar. Cada ação também precisa conseguir fazer retry:
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)
}
O código rapidamente vai se tornar verboso demais. Podemos encapsular a chave, incluindo a lógica de codificação e outras configs, em uma struct 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 ...
}
Depois, implementar alguns helpers para acessá-los rapidamente:
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) {
/* ... */
}
Por fim, podemos simplificar o uso:
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"))
// ...
}
Isso já é muito melhor!
A struct Timeslot contém apenas a definição do time slot. Precisamos de outra struct para encapsular seus estados:
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, /* ... */),
}
}
Deve haver apenas um único SlotStates para cada Timeslot em um pod, compartilhado entre o manager, o keeper e os workers. Por isso, é melhor ter um lugar centralizado para inicializá-los e armazená-los:
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()]
}
Ufa! Foi bastante coisa!! Fico feliz que você ainda esteja aqui!
Migrar grandes volumes de dados em um ambiente real é muito mais complexo do que um script simples consegue lidar. Ao desenhar e implementar cuidadosamente o código de migração com particionamento, processamento paralelo, gerenciamento de estados, tolerância a falhas, idempotência e acompanhamento centralizado do progresso, conseguimos alcançar um processo de migração confiável, manter a integridade dos dados e minimizar o downtime.
E dormir tranquilo à noite também! 😋
Eu sou o Oliver Nguyen -- engenheiro de software na Connectly.ai. Gosto de aprender e de me tornar uma versão melhor de mim mesmo todos os dias. De vez em quando, crio novos projetos open source. Compartilho conhecimento e reflexões durante essa jornada.
O post também foi publicado em olivernguyen.io.
©️ 2025 Todos os direitos reservados. Connectly Inc. Desenvolvido com ❤️, globalmente.