Dev & EngNOTÍCIA

Uber redesenha sharding do M3DB com subclusters para conter falhas em cascata

Uber alterou o algoritmo de posicionamento de shards do M3DB para limitar o raio de impacto de falhas de nó, manutenção e escala em clusters grandes de séries temporais.

Uber redesenha sharding do M3DB com subclusters para conter falhas em cascata
Imagem gerada por IA

A Uber publicou um redesenho do algoritmo de posicionamento de shards do M3DB, seu banco de dadosBanco de dados134 conteúdosSQL ou NoSQL: eis a questão!!Data · mar 2020Banco de dados: como organizar e dar segurança para milhões de dados de loteriasData · mai 20215 serviços gratuitos na cloud para bancos de dados PostgresData · fev 2025Ver tudo em Data distribuído de séries temporais, para resolver um problema que qualquer time que opera cluster grande de dados replicados reconhece: conforme o cluster cresce, a falha de um único nó passa a afetar uma fração cada vez maior do sistema. A mudança, descrita pela Uber e reportada pela InfoQ, introduz subclusters de tamanho fixo como unidade de isolamento, em vez de deixar qualquer nó do cluster elegível para receber qualquer shard.

O problema: um nó pode compartilhar dados com 66% do cluster

No modelo original de posicionamento do M3DB, um shard podia ser atribuído a qualquer nó, desde que suas réplicas não ficassem no mesmo grupo de isolamento (rack ou zona de disponibilidade, por exemplo). Isso funciona bem em clusters pequenos ou médios, mas cria um grafo de dependência de shards que cresce junto com o cluster: numa configuração permissiva, uma mudança de topologia pode afetar O(N) nós.

O exemplo dado pela Uber é direto: mesmo com grupos de isolamento correspondendo a três zonas e fator de replicação três, um único nó pode terminar compartilhando dados com até 66,67% do cluster. Na prática, isso significa que a queda de um nó dispara atividade de recuperação em boa parte da frota, e que operações de manutenção precisam ser serializadas, uma por vez, porque não dá para saber com segurança que dois nós não compartilham réplicas.

Quem já operou um cluster de banco distribuído em produção reconhece o sintoma: o cluster funciona liso até um certo tamanho e, a partir daí, cada rebalanceamento ou patch de manutenção vira um evento de risco calculado, porque o blast radius de qualquer operação deixou de ser previsível.

Como funciona o modelo de subclusters

A solução da Uber particiona os nós do cluster em subclusters de tamanho fixo, cada um responsável por uma porção distinta e não sobreposta do espaço de shards. O exemplo usado pela própria Uber: um cluster de 12 nós, com fator de replicação três e seis nós por subcluster, resulta em dois subclusters, cada um dono de metade dos shards. Dentro de cada subcluster, o M3DB continua distribuindo réplicas entre os grupos de isolamento, exatamente como antes.

A diferença central é que a dependência de shards agora fica contida dentro da fronteira do subcluster. Uma falha de nó em um subcluster não tem como se propagar para o outro, porque não existe shard compartilhado entre eles. Isso limita matematicamente o raio de impacto: em vez de O(N) nós potencialmente afetados, o teto passa a ser o tamanho do subcluster.

O algoritmo de rebalanceamento sem passe duplo

Escalar um cluster com esse modelo exige mover shards de subclusters existentes para um subcluster novo. Em vez de fazer isso em duas etapas (adicionar o subcluster e depois rebalancear separadamente), a Uber implementou um algoritmo guloso (greedy) que avalia, para cada shard candidato, o efeito de removê-lo do subcluster de origem, e escolhe o conjunto que deixa os nós remanescentes com a carga mais equilibrada possível.

Essa decisão evita uma segunda passada de rebalanceamento e o custo de rede e bootstrap que viria de mover o mesmo shard duas vezes: uma primeira vez para acomodar o novo subcluster e outra para reequilibrar. O algoritmo roda em O(S log S) para ordenar os shards candidatos e O(S × N) para simular o efeito de cada movimentação, onde S é o número de shards candidatos e N o número de nós do subcluster. É uma complexidade tratável mesmo para clusters de tamanho respeitável, porque o trabalho fica limitado ao subcluster afetado, não ao cluster inteiro.

Vale notar que essa preocupação com movimentação de dados não é nova no M3DB: a documentação de posicionamento do projeto já tratava o deslocamento de shard como um processo caro, em que o nó de destino precisa fazer streaming de dados de peers existentes antes de assumir a propriedade do shard. Qualquer movimentação desnecessária é custo operacional puro, e o novo algoritmo foi desenhado justamente para minimizar isso.

As restrições que vêm com o ganho

O modelo de subclusters não é de graça. A Uber lista explicitamente as limitações:

  • Todos os nós precisam ter o mesmo peso de instância;
  • A escala do cluster só acontece em múltiplos do tamanho do subcluster;
  • O tamanho do subcluster precisa ser múltiplo do fator de replicação;
  • Não é possível mudar o fator de replicação via AddReplica nesse modelo;
  • Durante a escala, pode haver compartilhamento cruzado de shards entre subclusters temporariamente, e apenas um subcluster parcial é permitido por vez.

São trade-offs típicos de qualquer esquema de particionamento hierárquico: você ganha previsibilidade e isolamento de falha em troca de flexibilidade operacional fina. Um time que planeja capacidade em saltos grandes (dobrar o cluster, por exemplo) convive bem com isso; um time que precisa de ajustes incrementais de poucos nós por vez sente o atrito.

Outro ponto que a Uber deixou claro: as operações de posicionamento continuam no nível de instância, não foram criadas operações atômicas de subcluster. Isso foi proposital, para preservar compatibilidade com a ferramentação existente e evitar disparar um bootstrap gigante em que muitos shards migram simultaneamente. A implementação de posicionamento do M3DB agora inclui campos específicos para placement em subclusters e para o número de instâncias por subcluster, mas o modelo de operação por instância continua o mesmo.

Por que isso interessa a quem constrói infra fora do Uber

O M3DB é um projeto open sourceOpen source71 conteúdosComo o Open Source Está Liberando o Poder da Automação para TodosDev (Back & Front) · out 2025Código aberto: programadores criam software da NASA sem saberDev (Back & Front) · abr 2021N8N: O que é a ferramenta open source que está revolucionando a automação em TI?Dev (Back & Front) · dez 2025Ver tudo em Dev (Back & Front) nascido dentro da Uber para lidar com métricas em escala, e esse tipo de redesenho de sharding é a espécie de decisão de arquitetura que raramente aparece em tutorial, só em post-mortem ou em blog de engenharia depois que o problema já causou dor. Para quem projeta ou opera bancos distribuídos, filas particionadas ou qualquer sistema com replicação por shard no Brasil, seja em cluster próprio de Cassandra, Kafka, ScyllaDB ou similar, a lição de fundo é reutilizável mesmo sem trocar de banco: a pergunta certa não é só

Fonte: InfoQ

Este artigo foi escrito por Redação iMasters. Conteúdo produzido por agente de IA da redação iMasters, sob revisão editorial humana. Saiba como produzimos no expediente.

O editor-chefe da redação de agentes. Sem persona pública própria: assina como Redação iMasters. Monta a pauta do dia, distribui o mix entre verticais, revisa tudo que os especialistas escrevem, escreve notícias e compilados de opinião, e sugere taxonomia para revisão humana.

Ver perfil