Connectly
Engenharia2024-10-22

Implementando uma máquina de estados distribuída com Redis para migrar bilhões de registros

Por Oliver Nguyen

Implementando uma máquina de estados distribuída com Redis para migrar bilhões de registros

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

Schema

Aqui está o schema simplificado da tabela room_events no Postgres:

room_events

  • id : ULID , chave primária
  • room_id : UUID , referencia rooms.id
  • data : JSON

Observações:

  • A coluna id é um ULID, que tem um componente de tempo.
  • Ela tem bilhões de registros.

Uma abordagem ingênua

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:

  • Escanear com paginação por cursor baseada na coluna id.
  • Para cada gravação, usar um batch de no máximo 25 registros (o número máximo de itens permitido pelo BatchWriteItem).
  • Escanear de trás para frente, da data do double-write até a data mais antiga. Assim temos os registros mais recentes disponíveis primeiro.
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(/* ... */ ) { /* ... */ }

Abordagem para o mundo real

Este script ingênuo vai operar sob algumas premissas:

  • As chamadas para queryRoomEvents() e batchWriteDynamo() sempre têm sucesso.
  • O número de registros é pequeno o suficiente para completar em uma única execução em uma única máquina.

Porém, no mundo real, não podemos confiar nessas premissas:

  • As requisições podem dar timeout ou falhar aleatoriamente.
  • Os bancos de dados podem atingir seus limites de throughput ou de conexão.
  • Os pods podem ser reiniciados por falhas ou novos deploys.
  • Um volume enorme de registros pode tornar a abordagem single-thread impraticável, levando uma eternidade para completar.

Para lidar com esses desafios, precisamos melhorar nosso script e implementar uma solução distribuída mais robusta:

Divisão e particionamento dos dados:

  • Dividir os dados agrupando os room events em Timeslots de 1 hora. Cada slot deve representar uma unidade de trabalho gerenciável que pode ser processada de forma independente. Isso permite melhor balanceamento de carga e maior paralelismo.

Processamento paralelo com múltiplos workers:

  • Rodar múltiplos workers em múltiplos pods, cada pod gerenciando múltiplas goroutines.
  • Cada worker vai pegar um Timeslot, migrá-lo, e depois passar para o próximo slot disponível, garantindo que todos os workers operem de forma concorrente.

Gerenciamento de estados no Redis:

  • Os workers vão adquirir um lock para cada Timeslot no Redis antes de processá-lo, garantindo que apenas um único worker possa tratar um slot por vez.
  • Depois de processar um Timeslot, o worker vai atualizar o status do slot para FINISHED.
  • Se um worker falhar, o lock será liberado ou expirado, permitindo que outro worker assuma e continue o processamento.

Keeper para acompanhamento centralizado do progresso:

  • Introduzir um Keeper para monitorar o progresso de todos os workers, manter uma visão geral global da migração e registrar o progresso global.
  • Ele vai acompanhar o progresso e o status de cada Timeslot, garantindo que nenhum slot fique de fora.

Tratamento de timeout e mecanismo de retry:

  • Implementar um mecanismo de retry com backoff exponencial tanto para leituras do PostgreSQL quanto para gravações no DynamoDB, tornando o script resiliente a falhas de rede transitórias ou timeouts.
  • E outra camada de retry e tratamento de erros para evitar qualquer pânico ou falha inesperada, garantindo que o script consiga se recuperar de forma graciosa, sem perda de dados ou estados corrompidos.

Verificação de dados:

  • Depois da migração, executar um script de verificação para garantir que todos os room events foram migrados com sucesso, sem registros ausentes.
  • O processo de verificação deve comparar a contagem de registros e amostras de dados entre PostgreSQL e DynamoDB.

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:

  • Se forem encontrados bugs no script e uma nova versão precisar ser executada, podemos resetar todos os estados e começar do zero sem excluir os dados já existentes no DynamoDB.
  • Se o script parar por qualquer motivo e retomar depois, ele pode sobrescrever alguns registros no DynamoDB. Isso é aceitável, já que regravações ocasionais são inofensivas e evitam a necessidade de refazer grandes partes do trabalho já concluído.

Arquitetura

Com esses conceitos em mente, vamos começar a pensar na arquitetura.

Armazenando estados no Redis

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:

  • Prefixar todas as chaves com mgre: → para conseguir escanear e apagar todas elas.
  • mgre:status → o status global da migração: NOT_STARTED, TO_START, IN_PROGRESS, TO_STOP, STOPPED, FINISHED.
  • mgre:{TIMESLOT}:worker → o worker atual trabalhando no slot, um lock exclusivo com TTL curto (por exemplo, 15-30 seg). Quando um worker falha ao liberar o lock, ele expira automaticamente para que outro worker possa assumir depois.
  • mgre:{TIMESLOT}:states → um JSON armazenando os estados do slot, com TTL longo (por exemplo, 1 semana) e propriedades:
  • .status → vazio, IN_PROGRESS ou FINISHED. Note que não existe status ERROR. Todos os room events e todos os slots precisam ser migrados com sucesso! 😎
  • .last_id → o último id de room events migrado, para conseguir retomar o progresso.

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.

Manager, Keepers e Workers

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:

  • Cada pod tem um único Manager.
  • Ele percorre todos os Timeslots.
  • Para cada slot, verifica o status e o lock correspondentes no Redis, checando mgre:{TIMESLOT}:states e mgre:{TIMESLOT}:worker.
  • Depois, distribui o slot para os workers, enviando-o para um channel do Go, slotCh.

Um Worker migra os room events de cada Timeslot:

  • Existem múltiplos workers, cada um rodando em uma goroutine separada.
  • Um worker recebe Timeslots do channel slotCh, um por vez.
  • Ele verifica o status em mgre:{TIMESLOT}:states novamente.
  • E adquire o lock com um comando SET NX em mgre:{TIMESLOT}:worker.
  • Se tudo tiver sucesso, ele vai começar a trabalhar no slot.
  • Ele vai salvar o progresso ocasionalmente, atualizando mgre:{TIMESLOT}:states.
  • E deve conseguir retomar a partir do progresso salvo.
  • Se falhar, o lock eventualmente expira e outro worker pode assumir.

Um Keeper monitora todo o progresso da migração:

  • Deve haver apenas um Keeper ativo entre todos os pods.
  • Devemos ter um lock mgre:keeper com TTL curto, então, se um Keeper de um pod falhar, outro Keeper pode assumir.
  • Um Keeper que falhou volta para o estado de standby e continua monitorando o lock para ver se pode se tornar ativo novamente.
  • A cada poucos segundos, o Keeper ativo escaneia a partir de mgre:last_slot para encontrar Timeslots FINISHED consecutivos. Depois, marca mgre:last_slot no slot FINISHED consecutivo mais recente e limpa esses slots.
  • Quando todos os slots restantes estiverem FINISHED, ele marca o mgre:status global como FINISHED. Todos os managers observam esse status e então param.
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 retry

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:

  • Sempre ter um recover(), porque qualquer pânico sem recover em uma goroutine pode parar todo o processo.
  • Existem 2 loops principais para conseguir reiniciar sempre que houver um problema:
  • O loop externo é responsável por reiniciar o loop interno. Ele só contém instruções simples, para garantir que nunca entre em pânico.
  • O loop interno é responsável por tratar a lógica: carregar estados, consultar dados, etc.
  • Para o Worker, há um loop ainda mais externo, para receber os slots e passá-los para o loop externo, que então passa para o loop interno para serem executados.
  • E uma função de retry para executar e tentar novamente cada consulta algumas vezes.
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:

  • A função de retry vai tentar novamente algumas vezes.
  • Se falhar, o loop interno vai parar e devolver o controle ao loop externo.
  • O loop externo então aguarda alguns segundos e reinicia o loop interno.
  • O loop interno então carrega os estados anteriores, o progresso, adquire o lock e continua.

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.

Outras observações

API para controlar a migração:

  • Devemos expor uma API para controlar a migração: iniciar e parar todos os managers, keepers e workers de todos os pods.
  • Ao receber o comando START, a API define mgre:status como TO_START ou IN_PROGRESS.
  • Quando os managers de todos os pods veem o status, eles iniciam ou retomam a migração.

Parar tudo usando context.WithCancel():

  • Usar um único context.Context e passá-lo para todos os managers, keepers e workers do mesmo pod.
  • Quando os managers veem o status mudar para TO_STOP, STOPPED, eles cancelam o context, fazendo o keeper e todos os workers pararem de forma graciosa.

Salvar o progresso como estados e renovar o lock periodicamente:

  • Enquanto cada worker ou keeper está trabalhando na sua tarefa, ele deve salvar o progresso no Redis e renovar o lock, para conseguir retomar depois e impedir que outros workers assumam o slot.

O Manager verifica periodicamente o mgre:last_slot no Redis:

  • Mesmo que o slot já tenha sido pego por um worker, há a chance de o worker parar e o slot não estar FINISHED.
  • Enquanto cada Manager percorre todos os slots, depois de um tempo ele precisa reiniciar o loop e verificar o status de mgre:last_slot no Redis, retomando a partir daí. Isso garante que não deixaremos nenhum slot inacabado.

Registrar e reportar o progresso:

  • O Keeper é responsável por monitorar o progresso, limpar os registros finalizados e atualizar mgre:last_slot.
  • Durante o job, o Keeper deve registrar e reportar o progresso periodicamente, para dar visibilidade ao status da migração e permitir o acompanhamento do progresso geral em tempo real.

Monitorar os recursos e limitar a taxa:

  • A migração pode consumir todos os recursos do banco de dados. Talvez seja necessário monitorar e adicionar limites de taxa, ajustando os parâmetros quando necessário, para garantir que tudo corra da forma mais tranquila possível.

Implementação

Timeslot

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

Retry() nas consultas

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

States

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!

SlotStates

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

Conclusão

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

Autor

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.