Sharding e Particionamento
Dividir os dados de um banco de dados em partições menores distribuídas entre múltiplos nós para escalar horizontalmente o armazenamento e a escrita — ao custo de complexidade operacional significativa e restrições nos padrões de acesso.
Intenção
Sharding divide os dados de um banco em partições menores (shards) espalhadas entre múltiplos nós, permitindo que o volume de dados e a taxa de escrita cresçam além do limite de uma única máquina.
Um único servidor de banco de dados tem limites físicos: capacidade de disco, memória RAM e largura de banda de I/O. Quando o volume de dados ultrapassa esses limites, escalar verticalmente (mais RAM, mais CPU, mais disco) se torna progressivamente mais caro e eventualmente impossível. Sharding distribui a carga por múltiplos nós — cada um responsável por uma fatia do dataset — transformando o problema de capacidade vertical em um problema de coordenação horizontal.
Problema
Sistemas que crescem além de dezenas ou centenas de gigabytes de dados ativos enfrentam gargalos que indexação e otimização de queries não conseguem resolver sozinhos:
- Volume de dados excede a capacidade de um nó: tabelas com bilhões de linhas degradam a performance mesmo com índices otimizados. O índice em si se torna grande demais para caber em memória, forçando I/O de disco em cada operação de busca.
- Taxa de escrita excede o throughput de um nó: um único servidor tem um número máximo de writes por segundo que pode processar. Em sistemas de alto volume — logs, métricas, transações financeiras — esse limite é atingido antes do limite de armazenamento.
- Read replicas não resolvem o problema de escrita: adicionar replicas de leitura distribui a carga de leitura, mas todas as escritas ainda passam pelo primary. O gargalo de escrita permanece.
Como funciona
Particionamento horizontal vs vertical
É importante distinguir os dois tipos de particionamento, pois o termo frequentemente causa confusão:
- Particionamento horizontal (sharding): divide as linhas de uma tabela entre múltiplos nós. Todos os shards têm o mesmo schema. Ex.: usuários com ID 1–1.000.000 no Shard A, ID 1.000.001–2.000.000 no Shard B. É o que se entende por sharding na maioria dos contextos.
-
Particionamento vertical: divide as colunas de uma
tabela entre bancos diferentes. Ex.: a tabela
userscom dados de perfil em um banco e dados de pagamento em outro. Técnica diferente, com motivação diferente — normalmente usada para isolar domínios por segurança ou compliance, não para escala de volume.
Arquitetura de sharding
O componente central de uma arquitetura shardada é o shard router: a camada que recebe a requisição do cliente, determina em qual shard os dados residem e encaminha a operação ao nó correto.
Cliente
│
▼
┌──────────────┐
│ Shard Router │ determina o shard com base na shard key
└──────┬───────┘
│
┌────┼────┐
▼ ▼ ▼
┌───┐┌───┐┌───┐
│ A ││ B ││ C │ cada shard é um nó de banco independente
└───┘└───┘└───┘
(mesmo schema em todos os shards)
O router pode ser implementado na camada de aplicação (o código decide qual banco conectar), em um proxy dedicado (ex.: ProxySQL, Vitess) ou ser nativo ao banco (ex.: MongoDB com mongos, Cassandra com token-aware drivers).
Shard key
A shard key é o campo (ou combinação de campos) usado para determinar em qual shard um registro pertence. É a decisão de design mais crítica em uma arquitetura shardada — uma shard key mal escolhida cria hotspots: um shard sobrecarregado enquanto os demais ficam ociosos, anulando o benefício da distribuição.
Critérios para uma boa shard key:
- Alta cardinalidade: muitos valores distintos possíveis, para que os dados se distribuam entre muitos shards.
- Distribuição uniforme: os valores não devem se concentrar em poucos ranges ou hashes.
- Alinhamento com os padrões de acesso: operações frequentes devem tocar um shard por vez. Evite shard keys que forcem o router a consultar todos os shards para responder uma query comum.
Estratégias de sharding
Range-based
O shard é determinado por um intervalo de valores da shard key. Ex.:
registros com created_at em janeiro vão ao Shard A, fevereiro
ao Shard B.
shard key: user_id (numérico)
Shard A: user_id 1 – 1.000.000
Shard B: user_id 1.000.001 – 2.000.000
Shard C: user_id 2.000.001 – 3.000.000
Vantagem: range queries eficientes ("todos os users de 1 a 500k")
Risco: hotspot se novos registros são criados sequencialmente
(todos vão ao último shard enquanto os anteriores ficam ociosos)
Hash-based
O shard é determinado por uma função de hash aplicada à shard key:
shard = hash(shard_key) % num_shards. Distribui os dados
uniformemente, eliminando hotspots de range.
shard key: user_id
shard = MD5(user_id) % 3
user_id=1001 → hash mod 3 = 2 → Shard C
user_id=1002 → hash mod 3 = 0 → Shard A
user_id=1003 → hash mod 3 = 1 → Shard B
Vantagem: distribuição uniforme, sem hotspots de range
Risco: range queries ineficientes (exige consultar todos os shards)
resharding exige remapear quase todos os registros
Directory-based
Uma tabela de lookup (o "diretório") mapeia explicitamente cada valor de shard key para seu shard. O router consulta o diretório antes de cada operação.
Tabela de lookup (diretório):
┌──────────┬────────┐
│ tenant │ shard │
├──────────┼────────┤
│ empresa1 │ A │
│ empresa2 │ B │
│ empresa3 │ A │
│ empresa4 │ C │
└──────────┴────────┘
Vantagem: flexível — pode mover tenants entre shards sem alterar lógica
Risco: a tabela de lookup vira ponto único de falha e gargalo de latência
Resharding e consistent hashing
Quando um shard fica grande demais ou sobrecarregado, os dados precisam ser
redistribuídos — operação chamada de resharding. Com
hash-based simples (% N), adicionar um shard muda a fórmula e
potencialmente invalida o mapeamento de quase todos os registros, exigindo
mover a maioria dos dados.
Consistent hashing resolve isso: os shards e as chaves são
posicionados num anel circular, e cada chave pertence ao shard mais próximo
no sentido horário. Ao adicionar ou remover um shard, apenas as chaves
adjacentes a ele no anel precisam ser movidas — em média K/N
chaves (onde K é o total de chaves e N o número de shards), em vez de quase
todas.
Joins cross-shard
Joins entre dados que residem em shards diferentes não são executados pelo banco — cada shard é um banco independente. A aplicação precisa buscar os dados de cada shard separadamente e combinar os resultados em memória. Dependendo do volume, isso pode ser mais lento que o join original e introduz lógica de merge na camada de aplicação.
A solução mais comum é desnormalizar: duplicar os dados necessários dentro do mesmo shard para evitar o join cross-shard. Ex.: em vez de fazer join com a tabela de usuários em outro shard, copiar o nome do usuário para dentro do registro do pedido no momento da criação.
Quando usar
- Volume de dados excede a capacidade de um único nó: quando escalar verticalmente (mais RAM, mais disco) se torna proibitivo em custo ou inviável tecnicamente.
- Taxa de escrita excede o throughput de um nó: sistemas de alto volume como analytics em tempo real, séries temporais de IoT ou plataformas de e-commerce em grande escala.
- Os dados têm uma dimensão natural de particionamento: por usuário (redes sociais, SaaS multi-tenant), por região geográfica ou por período de tempo. A shard key surge naturalmente do domínio.
Quando evitar
- Esgote as alternativas mais simples primeiro: indexação correta, otimização de queries, caching e read replicas podem adiar ou eliminar a necessidade de sharding por muito tempo. Sharding adiciona complexidade operacional permanente.
- Joins cross-shard são frequentes no seu workload: se o padrão de acesso do sistema exige combinar dados de múltiplas dimensões com frequência, sharding transforma cada operação simples em um problema de coordenação distribuída.
- Sem uma shard key natural: se não há uma chave que distribua os dados uniformemente e que se alinhe com os padrões de acesso, o sistema vai criar hotspots. Sharding sem uma boa shard key é pior que não shardear.
Prós e contras
Prós
- Escala horizontal real de escrita: a capacidade de ingestão de dados cresce linearmente com o número de shards.
- Volume de dados ilimitado: o dataset pode crescer além do limite de qualquer máquina individual.
- Isolamento de falhas: a falha de um shard afeta apenas os dados naquele shard, não o sistema inteiro.
- Queries localizadas ficam mais rápidas: com uma boa shard key, a maioria das queries toca apenas um shard — um dataset menor, com índices menores que cabem inteiros em memória.
Contras
- Complexidade operacional alta: monitorar, fazer backup, restaurar e atualizar N bancos de dados independentes em vez de um.
- Joins cross-shard: precisam ser implementados na aplicação, com lógica de merge manual e múltiplas viagens de rede.
- Transações cross-shard: ACID não se aplica entre shards. Two-phase commit ou padrões como Saga são necessários, com trade-offs significativos de consistência.
- Resharding é custoso: redistribuir dados com o sistema em produção exige coordenação cuidadosa para evitar inconsistência e downtime.
- Hotspots são difíceis de diagnosticar: a distribuição pode parecer uniforme nos dados mas ser extremamente desigual no acesso.
Armadilhas comuns
1. Shard key que cria hotspot
O exemplo clássico é usar created_at (timestamp) como shard key
num sistema de logs ou eventos. Todo o tráfego novo vai ao shard do período
mais recente — os shards antigos ficam ociosos enquanto o shard atual é
bombardeado. O sistema shardado performa pior que o sistema não shardado
porque a carga, em vez de distribuída, está concentrada num único nó.
Regra prática: antes de adotar uma shard key, simule a distribuição com dados reais de produção (não com dados de teste uniformes) e com a projeção do padrão de acesso futuro, não apenas do presente.
2. Joins cross-shard em loop
A aplicação busca uma lista de IDs do Shard A e, para cada ID, faz uma query no Shard B para buscar dados relacionados — o problema N+1 em escala distribuída. Em vez de um join eficiente no banco, são N queries de rede serializadas. O resultado é pior em latência e carga do que antes do sharding.
Soluções: desnormalizar os dados necessários para dentro do mesmo shard no momento da escrita, ou usar fan-out (consultar todos os shards em paralelo e agregar os resultados), aceitando o custo de múltiplas conexões.
3. Resharding não planejado
O crescimento do sistema excede a capacidade dos shards existentes e é necessário adicionar shards. Com hash-based simples, isso exige mover a maioria dos dados. Com o sistema em produção, a janela de consistência durante a migração precisa ser gerenciada — dados sendo movidos não podem ser acessados pelo router antigo e o router novo ainda não conhece o destino final. Planejar resharding antes de precisar dele é parte da arquitetura.
4. Transações cross-shard
Uma operação que precisa modificar dados em dois shards diferentes não tem garantia ACID. Se o sistema confirma a escrita no Shard A mas falha antes de escrever no Shard B, os dados ficam em estado inconsistente. Two-phase commit (2PC) resolve o problema mas introduz locks distribuídos e degrada a disponibilidade. O padrão Saga substitui a transação distribuída por uma sequência de transações locais com compensações em caso de falha.
Arquiteturas e padrões relacionados
Replicação e sharding são técnicas complementares e frequentemente usadas juntas: sharding distribui as escritas entre múltiplos nós primários, enquanto replicação adiciona replicas de leitura a cada shard para aumentar o throughput de leitura e a disponibilidade. A arquitetura resultante tem tanto escala horizontal de escrita quanto tolerância a falhas.
Bancos NoSQL como Cassandra, DynamoDB e MongoDB foram projetados com sharding nativo — a lógica de distribuição faz parte do banco, não da aplicação. Em bancos SQL relacionais, sharding normalmente é implementado na camada de aplicação ou com ferramentas como Vitess (MySQL) ou Citus (PostgreSQL). A escolha entre SQL e NoSQL é influenciada por quão central o sharding é para o design desde o início.
O caching é frequentemente a primeira linha de defesa contra sobrecarga do banco: ao reduzir o número de leituras que chegam ao banco, o caching adia o momento em que sharding se torna necessário. Avaliar caching antes de sharding é a sequência correta de decisões arquiteturais.