Por Oliver Nguyen

Como empresa startup, empezamos con una arquitectura simple: todos los datos se almacenaban en una única instancia centralizada de Postgres, compartida por unos pocos servicios. Ha funcionado bien y nos ha permitido movernos rápido, entregar funcionalidades, entrar al mercado, adquirir clientes y crecer exponencialmente.
Avanzando un par de años, hoy tenemos miles de clientes y miles de millones de registros; es momento de migrar las tablas principales de Postgres a una mejor solución. Elegimos DynamoDB, un almacén de datos clave-valor de AWS, con alta disponibilidad, rendimiento de baja latencia, y vemos las consultas reducirse a milisegundos de un solo dígito.
Como en cualquier buen plan de migración, analizamos cuidadosamente los usos, definimos el esquema y los índices, comenzamos con doble escritura, cambiamos las consultas de lectura a DynamoDB con respaldo a Postgres, importamos todos los registros a DynamoDB, y finalmente dejamos de consultar desde Postgres.
Este artículo se enfoca en el paso de importar registros de la tabla room_events: escanear todos los registros de esta tabla y escribirlos en DynamoDB*. El objetivo es asegurar una migración completa, sin eventos perdidos. El script de migración debe ser rápido, capaz de ejecutarse en paralelo, y capaz de detenerse y reanudarse desde su último estado. También debe ser resiliente, manejando errores de red o interrupciones y reanudando sin problemas.*
¡Veamos cómo podemos hacer esto! 🥸
Aquí está el esquema simplificado de la tabla room_events en Postgres:
room_events
Notas:
Escribamos una versión simple del script de migración: ejecutar en un solo bucle, escanear todos los registros, y escribir en DynamoDB. Además, incluyamos algunas cosas triviales para no tener que preocuparnos por ellas más adelante:
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 ingenuo operará bajo algunas suposiciones:
Sin embargo, en el mundo real, no podemos confiar en esas suposiciones:
Para manejar estos desafíos, necesitamos mejorar nuestro script e implementar una solución distribuida más robusta:
División y particionamiento de datos:
Procesamiento paralelo con múltiples workers:
Gestión de estados en Redis:
Keeper para seguimiento centralizado del progreso:
Manejo de timeouts y mecanismo de reintento:
Verificación de datos:
Aprovechando la idempotencia en las escrituras:
Otro punto importante es que el paso de escritura para cada registro es idempotente. Esto significa que escribir el mismo registro en DynamoDB múltiples veces simplemente sobrescribirá la versión anterior, asegurando que no haya duplicados.
Podemos usar esto a nuestro favor:
Con esos conceptos en mente, empecemos a pensar en la arquitectura.
Dado que ya dividimos los registros en Timeslots, cada uno consiste en todos los eventos de sala de una hora. Necesitamos hacer seguimiento del progreso de estas franjas en Redis:
Así que necesitamos 2 claves para cada franja. Multiplicado por 3 años de datos, habrá 52.560 claves (3 * 365 * 24 * 2). ¡Eso es mucho!
Podemos optimizar las claves de Redis limpiando el grupo de franjas FINISHED consecutivas y reemplazándolas con una única clave mgre:last_slot. De esta manera, solo necesitamos mantener el seguimiento de una pequeña cantidad de claves. Por supuesto, esto se basa en la suposición de que las franjas tienen un número similar de eventos de sala, para que cada worker pueda tomar una cantidad de tiempo similar en terminar.
El script se ejecuta en múltiples pods. Cada pod consiste en un Manager, un Keeper, y múltiples Workers. Cada uno se ejecuta en su propia goroutine.
Un Manager gestiona todos los workers en un pod:
Un Worker migra eventos de sala por cada Timeslot:
Un Keeper supervisa todo el progreso de la migración:
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)
}
}
}
El plan se ve bien, ¿verdad? No, todavía no. ¿Qué pasa si alguno de los pasos anteriores falla?
Cada Manager, Worker, Keeper corre en una goroutine y siempre se reintenta a sí mismo:
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
}
Cuando ocurre cualquier error, por ejemplo, un timeout de red:
Esto asegura que el código siempre se ejecute hasta que todos los registros se migren, o hasta que el pod se reinicie. En el último caso, reanudará el progreso de la migración la próxima vez.
API para controlar la migración:
Detener todo usando context.WithCancel():
Guardar el progreso como estados y refrescar el lock periódicamente:
El Manager verifica periódicamente el mgre:last_slot en Redis:
Registrar y reportar el progreso:
Monitorear los recursos y el límite de tasa:
Cada Timeslot representa todos los eventos de sala en una hora. Podemos implementarlo como un time.Time y almacenarlo en Redis como un string con el 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 discutimos antes, tenemos 2 bucles para manejar panic, reintento y reanudación de estados. Así que para esta función retry(), solo necesitamos reintentar la consulta un par de veces, para poder tolerar algunas fallas de red:
// 👉 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 una lógica de reintento simple:
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
}
}
}
Un ejemplo de implementación de estrategia de reintento, donde cada reintento ocurre después de una duración predefinida:
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
}
}
Y la implementación de la función msgf(), que se puede usar para devolver rápidamente un 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...)
}
}
Hay muchos estados: estado global, progreso global, estado de franja, progreso de franja, locks, etc. Para cada estado, necesitamos refrescar, actualizar o eliminar. Cada acción también necesita poder reintentarse:
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)
}
El código rápidamente se vuelve demasiado verboso. Podemos encapsular la clave, incluyendo la lógica de codificación y otras configuraciones, en un 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 ...
}
Luego implementamos algunos helpers para acceder a ellos rápidamente:
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) {
/* ... */
}
Finalmente, podemos simplificar el 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"))
// ...
}
¡Eso es mucho mejor!
El struct Timeslot solo contiene la definición para la franja de tiempo. Necesitamos otro struct para encapsular sus 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, /* ... */),
}
}
Debería haber solo un único SlotStates para cada Timeslot en un pod, compartido entre el manager, el keeper y los workers. Así que es mejor tener un lugar centralizado para inicializarlos y almacenarlos:
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()]
}
¡Uf! ¡Eso fue mucho! Me alegra que sigas aquí!
Migrar grandes volúmenes de datos en un entorno real es mucho más complejo de lo que un script simple puede manejar. Al diseñar e implementar cuidadosamente el código de migración con particionamiento, procesamiento paralelo, gestión de estados, tolerancia a fallos, idempotencia y seguimiento centralizado del progreso, podemos lograr un proceso de migración confiable, mantener la integridad de los datos y minimizar el downtime.
¡Y también dormir bien por la noche! 😋
Soy Oliver Nguyen -- ingeniero de software en Connectly.ai. Disfruto aprender y ver una mejor versión de mí mismo cada día. Ocasionalmente lanzo nuevos proyectos de código abierto. Comparto conocimiento y reflexiones durante mi camino.
Esta publicación también está publicada en olivernguyen.io.
©️ 2025 Todos los derechos reservados. Connectly Inc. Creado con ❤️, globalmente.