Connectly
Ingeniería2024-10-22

Implementa una máquina de estados distribuida con Redis para migrar miles de millones de registros

Por Oliver Nguyen

Implementa una máquina de estados distribuida con Redis para migrar miles de millones de registros

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! 🥸

Esquema

Aquí está el esquema simplificado de la tabla room_events en Postgres:

room_events

  • id : ULID , clave primaria
  • room_id : UUID , referencia a rooms.id
  • data : JSON

Notas:

  • La columna id es ULID, que tiene un componente temporal.
  • Tiene miles de millones de registros.

Un enfoque ingenuo

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:

  • Escanear con paginación por cursor basada en la columna id.
  • Para cada escritura, usar un lote de máximo 25 registros (el número máximo de elementos permitido por BatchWriteItem).
  • Escanear hacia atrás desde la fecha de doble escritura hasta la fecha más temprana. Así tenemos los registros más recientes disponibles primero.
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(/* ... */ ) { /* ... */ }

El enfoque del mundo real

Este script ingenuo operará bajo algunas suposiciones:

  • Las llamadas a queryRoomEvents() y batchWriteDynamo() siempre tienen éxito.
  • El número de registros es suficientemente pequeño para completar en una sola ejecución en una sola máquina.

Sin embargo, en el mundo real, no podemos confiar en esas suposiciones:

  • Las solicitudes pueden agotar el tiempo o fallar aleatoriamente.
  • Las bases de datos pueden alcanzar sus límites de throughput o de conexiones.
  • Los pods pueden reiniciarse debido a fallas o nuevos despliegues.
  • Una gran cantidad de registros puede hacer que el enfoque de un solo hilo sea poco práctico, tomando una eternidad en completarse.

Para manejar estos desafíos, necesitamos mejorar nuestro script e implementar una solución distribuida más robusta:

División y particionamiento de datos:

  • Dividir los datos agrupando los eventos de sala en Timeslots de 1 hora. Cada franja debería representar una unidad de trabajo manejable que se procese de forma independiente. Esto permite un mejor balanceo de carga y un paralelismo mejorado.

Procesamiento paralelo con múltiples workers:

  • Ejecutar múltiples workers en múltiples pods, con cada pod gestionando múltiples goroutines.
  • Cada worker tomará un Timeslot, lo migrará, y luego avanzará a la siguiente franja disponible, asegurando que todos los workers operen concurrentemente.

Gestión de estados en Redis:

  • Los workers adquirirán un lock para cada Timeslot en Redis antes de procesarlo, para asegurar que solo un worker pueda manejar una franja a la vez.
  • Después de procesar un Timeslot, el worker actualizará el estado de la franja a FINISHED.
  • Si un worker falla, el lock se liberará o expirará, permitiendo que otro worker lo tome y continúe procesando.

Keeper para seguimiento centralizado del progreso:

  • Introducir un Keeper para monitorear el progreso de todos los workers, mantener una vista global de la migración, y registrar el progreso global.
  • Hará seguimiento del progreso y estado de cada Timeslot, asegurando que no se pierda ninguna franja.

Manejo de timeouts y mecanismo de reintento:

  • Implementar un mecanismo de reintento con backoff exponencial tanto para las lecturas de PostgreSQL como para las escrituras de DynamoDB, haciendo que el script sea resiliente a fallas de red transitorias o timeouts.
  • Y otra capa de reintento y manejo de errores para prevenir cualquier panic o falla inesperada, asegurando que el script pueda recuperarse elegantemente, sin pérdida de datos ni estados corruptos.

Verificación de datos:

  • Después de la migración, ejecutar un script de verificación para asegurar que todos los eventos de sala se hayan migrado exitosamente, sin registros faltantes.
  • El proceso de verificación debería comparar el conteo de registros y datos de muestra entre PostgreSQL y DynamoDB.

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:

  • Si se encuentran bugs en el script y hay que ejecutar una nueva versión, podemos reiniciar todos los estados y empezar desde cero sin eliminar los datos existentes en DynamoDB.
  • Si el script se detiene por cualquier razón y se reanuda más tarde, podría sobrescribir algunos registros en DynamoDB. Esto es aceptable, ya que las re-escrituras ocasionales son inofensivas y evitan la necesidad de rehacer grandes porciones de trabajo ya completado.

Arquitectura

Con esos conceptos en mente, empecemos a pensar en la arquitectura.

Almacenar estados en Redis

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:

  • Prefijar todas las claves con mgre: → para poder escanear y eliminarlas todas.
  • mgre:status → el estado global de la migración: NOT_STARTED, TO_START, IN_PROGRESS, TO_STOP, STOPPED, FINISHED.
  • mgre:{TIMESLOT}:worker → el worker actual trabajando en la franja, un lock exclusivo con TTL corto (por ejemplo, 15-30 seg). Cuando un worker falla en liberar el lock, expirará automáticamente para que otro worker pueda tomarlo más adelante.
  • mgre:{TIMESLOT}:states → un JSON que almacena los estados de la franja, con TTL largo (por ejemplo, 1 semana) y propiedades:
  • .status → vacío, IN_PROGRESS o FINISHED. Nota que no hay estado ERROR. ¡Todos los eventos de sala y todas las franjas deben migrarse exitosamente! 😎
  • .last_id → el último id de evento de sala migrado, para poder reanudar el progreso.

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.

Manager, Keepers y Workers

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:

  • Cada pod tiene un único Manager.
  • Recorre todos los Timeslots.
  • Para cada franja, verifica el estado y el lock correspondientes en Redis revisando mgre:{TIMESLOT}:states y mgre:{TIMESLOT}:worker.
  • Luego despacha la franja a los workers enviándola a un canal de Go slotCh.

Un Worker migra eventos de sala por cada Timeslot:

  • Hay múltiples workers, cada uno ejecutándose en una goroutine separada.
  • Un worker recibe Timeslots del canal slotCh, uno a la vez.
  • Verifica el estado en mgre:{TIMESLOT}:states de nuevo.
  • Y adquiere el lock mediante un comando SET NX en mgre:{TIMESLOT}:worker.
  • Si todo tiene éxito, empezará a trabajar en la franja.
  • Guardará el progreso ocasionalmente actualizando mgre:{TIMESLOT}:states.
  • Y debería poder reanudar desde el progreso guardado.
  • Si falla, el lock eventualmente expirará y otro worker podrá tomarlo.

Un Keeper supervisa todo el progreso de la migración:

  • Debería haber solo un Keeper activo para todos los pods.
  • Deberíamos tener un lock mgre:keeper con TTL corto, así que si un Keeper de un pod falla, otro Keeper puede tomar el control.
  • Un Keeper fallido volverá al estado de espera, y continuará monitoreando el lock para ver si puede volverse activo de nuevo.
  • Cada pocos segundos, el Keeper activo escanea desde mgre:last_slot para encontrar Timeslots FINISHED consecutivos. Luego marcará mgre:last_slot en la última franja FINISHED consecutiva, y las limpiará.
  • Cuando todas las franjas restantes estén FINISHED, marca el mgre:status global como FINISHED. Todos los managers observan ese estado y luego se detendrá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)
        }
    }
}

Mecanismo de reintento

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:

  • Siempre tener recover(), porque cualquier panic sin recover en una goroutine puede detener todo el proceso.
  • Hay 2 bucles principales para poder reiniciar cuando hay un problema:
  • El bucle externo es responsable de reiniciar el bucle interno. Solo contiene declaraciones simples para asegurar que nunca haga panic.
  • El bucle interno es responsable de manejar la lógica: cargar estados, consultar datos, etc.
  • Para el Worker, hay un bucle extra más externo para recibir franjas y pasarlas al bucle externo y luego al bucle interno para su ejecución.
  • Y una función de reintento para ejecutar y reintentar cada consulta algunas veces más.
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:

  • La función de reintento reintentará algunas veces.
  • Si falla, el bucle interno se detendrá y devolverá el control al bucle externo.
  • El bucle externo ahora espera algunos segundos y reinicia el bucle interno.
  • El bucle interno entonces carga los estados anteriores, el progreso, adquiere el lock, y continúa.

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.

Otras notas

API para controlar la migración:

  • Deberíamos exponer una API para controlar la migración: iniciar y detener todos los managers, keepers, workers de todos los pods.
  • Al recibir el comando START, la API establece mgre:status en TO_START o IN_PROGRESS.
  • Cuando los managers de todos los pods ven el estado, iniciarán o reanudarán la migración.

Detener todo usando context.WithCancel():

  • Usar un único context.Context y pasarlo a todos los managers, keepers y workers en el mismo pod.
  • Cuando los managers ven que el estado cambia a TO_STOP, STOPPED, cancelarán el contexto, haciendo que el keeper y todos los workers se detengan elegantemente.

Guardar el progreso como estados y refrescar el lock periódicamente:

  • Mientras cada worker o keeper está trabajando en lo suyo, debería guardar el progreso en Redis, y refrescar el lock, para poder reanudar más adelante y evitar que otros workers tomen la franja.

El Manager verifica periódicamente el mgre:last_slot en Redis:

  • Aunque la franja ya esté tomada por un worker, existe la posibilidad de que el worker se detenga y la franja no esté FINISHED.
  • Mientras cada Manager recorre todas las franjas, después de un rato, necesita reiniciar el bucle y verificar el estado de mgre:last_slot en Redis, luego reanudar desde ahí. Esto es para asegurar que no dejemos ninguna franja sin terminar.

Registrar y reportar el progreso:

  • El Keeper es responsable de supervisar el progreso, limpiar los registros terminados, y actualizar mgre:last_slot.
  • Durante el trabajo, el Keeper debería registrar y reportar el progreso periódicamente para dar visibilidad sobre el estado de la migración y permitir el seguimiento en tiempo real del progreso general.

Monitorear los recursos y el límite de tasa:

  • La migración puede consumir todos los recursos de la base de datos. Es posible que queramos monitorear y agregar límites de tasa, ajustar finamente los parámetros cuando sea necesario, para asegurar que todo transcurra lo más fluido posible.

Implementación

Timeslot

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)}
}

Reintento de consultas con Retry()

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...)
   }
}

Estados

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!

SlotStates

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()]
}

Conclusión

¡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! 😋

Autor

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.