Configuração de Sharding por Hash
Para avaliar a distribuição de carga em um ambiente distribuído, o prmieiro passo consiste em habilitar o sharding no banco de dados alvo e definir a estratégia de fragmentação para a coleção específica. A abordagem por hash no campo _id garante uma distribuição uniforme dos documentos entre os shards.
// Habilitando sharding no banco de dados da aplicação
sh.enableSharding("app_db");
// Aplicando estratégia de hash no campo _id da coleção
sh.shardCollection("app_db.user_records", { _id: "hashed" });
Definição do Modelo de Carga (Workload)
O framework Yahoo! Cloud Serving Benchmark (YCSB) é amplamente utilizado para testes de performence em bancos NoSQL. Sua arquitetura permite a customização de parâmetros como proporção de leitura/escrita, tamanho dos registros e distribuição de probabilidade. Para esta avaliação, foi desenhado um perfil de carga que simula o comportamento real da aplicação em produção.
| Identificador | Descrição do Perfil | Volume de Registros |
|---|---|---|
| business_profile | 60% Leituras, 40% Inserções | 1.000.000 |
Métricas de Avaliação
Os indicadores primários monitorados durante os testes incluem:
- Tempo de Execução (RunTime): Duração total do teste.
- Vazão (Throughput): Número de operações processadas por segundo.
- Latência Média (Average Latency): Tempo de resposta médio por operação.
O ponto de inflexão (gargalo) é identificado ao incrementar o número de threads até que a vazão (ops/s) estabilize ou decresça, enquanto a latência média e a utilização de CPU nos nós do cluster atingem níveis críticos.
Configuração e Execução do Workload Personalizado
O script abaixo detalha a criação do arquivo de configuração do YCSB e a execução das fases de carga e transação, utilizando variáveis de ambiente para facilitar a iteração com diferentes níveis de concorrência.
# Criação do arquivo de configuração do perfil de carga
cat > profiles/business_profile.cfg << 'EOF'
recordcount=1000000
operationcount=1000000
workload=site.ycsb.workloads.CoreWorkload
readallfields=true
readproportion=0.6
updateproportion=0.0
scanproportion=0.0
insertproportion=0.4
requestdistribution=zipfian
EOF
CONCURRENCY_LEVEL=100
MONGO_URI="mongodb://admin:'secure_pass'@host1:20000,host2:20000,host3:20000/app_db?authSource=admin"
# Fase de carga de dados (Data Load)
bin/ycsb load mongodb -threads $CONCURRENCY_LEVEL -P profiles/business_profile.cfg \
-p fieldcount=1 -p fieldlength=1024 -p clientbuffering=true \
-p table=user_records -p mongodb.url=$MONGO_URI \
> "results/load_phase_${CONCURRENCY_LEVEL}.out" 2> "logs/load_phase.err" &
tail -f "results/load_phase_${CONCURRENCY_LEVEL}.out"
# Fase de execução de transações (Transaction Run)
bin/ycsb run mongodb -threads $CONCURRENCY_LEVEL -P profiles/business_profile.cfg \
-p fieldcount=1 -p fieldlength=1024 -p clientbuffering=true \
-p table=user_records -p mongodb.url=$MONGO_URI \
> "results/run_phase_${CONCURRENCY_LEVEL}.out" 2> "logs/run_phase.err" &
tail -f "results/run_phase_${CONCURRENCY_LEVEL}.out"
Análise Estatística dos Resultados com Sharding
A tabela a seguir apresenta os dados coletados ao variar o parâmetro CONCURRENCY_LEVEL no cluster com sharding habiliatdo.
| Threads | Tempo (s) | Vazão (ops/s) | Ops Leitura | Latência Leitura (µs) | Ops Inserção | Latência Inserção (µs) |
|---|---|---|---|---|---|---|
| 1 | 1199 | 833 | 599849 | 1151 | 400151 | 1262 |
| 10 | 145 | 6916 | 600565 | 1486 | 399435 | 1274 |
| 20 | 71 | 14097 | 598887 | 1357 | 401113 | 1377 |
| 30 | 52 | 19102 | 599719 | 1495 | 400281 | 1521 |
| 50 | 52 | 19126 | 600574 | 2833 | 399426 | 1835 |
| 80 | 69 | 14470 | 599662 | 4440 | 400338 | 2553 |
| 100 | 42 | 23628 | 599939 | 3981 | 400061 | 3827 |
| 150 | 45 | 21845 | 601056 | 5593 | 398944 | 4284 |
| 200 | 47 | 20879 | 599388 | 6972 | 400612 | 5275 |
| 300 | 81 | 12336 | 599795 | 10121 | 400205 | 8441 |
A análise de telemetria dos nós (16 vCPUs cada) indica que, com 20 threads, a ociosidade de CPU mantém-se entre 50% e 60%. Durante a fase de carga com 50 threads (100% inserções), a ociosidade cai para menos de 20%, indicando saturação devido ao overhead de gravação. Na fase de transação (60% leitura / 40% escrita) com 50 threads, a carga equilibra-se novamente com 60% de ociosidade.
Quando a concorrência atinge 200 a 300 threads, o nó primário do roteador (mongos) e os shards começam a apresentar contenção severa, com a ociosidade do nó 1 caindo para 5%. O ponto de equilíbrio ideal entre latência e vazão máxima é observado em 100 threads simultâneas.
Comparativo: Cenário sem Sharding
Para isolar o impacto da fragmentação de dados, o mesmo teste foi executado em uma arquitetura de replica set sem sharding.
| Threads | Tempo (s) | Vazão (ops/s) | Ops Leitura | Latência Leitura (µs) | Ops Inserção | Latência Inserção (µs) |
|---|---|---|---|---|---|---|
| 100 | 189 | 5279 | 600016 | 7094 | 399984 | 2572 |
| 100 | 38 | 26168 | 600205 | 4076 | 399795 | 2696 |
| 150 | 39 | 25371 | 600090 | 5650 | 399910 | 5884 |
| 150 | 85 | 11764 | 599122 | 6058 | 400878 | 3877 |
Neste cenário, a inspeção dos metadados do banco revela que todos os documentos foram alocados em um único shard (shard2), cujos dados residem fisicamente no node2 (primário) e node3 (secundário). A carga de CPU concentrou-se massivamente no nó primário, enquanto o nó secundário e o árbitro (node1) permaneceram subutilizados. Embora o cluster suporte picos de 150 threads sem degradação catastrófica, a ausência de sharding impede a escalabilidade horizontal linear.
Diretrizes de Arquitetura e Otimização
Os dados coletados demonstram que, operando entre 100 e 200 conexões simultâneas em um cluster com sharding, a arquitetura atinge facilmente uma vazão superior a 20.000 ops/s. Para um volume de 1 milhão de operações (60% leitura, 40% inserção), o tempo total de processamento fica em torno de 40 segundos, com latência média inferior a 10ms por requisição.
Para operações de carga ou processamento em massa (bulk operations), a execução single-thread resulta em taxas de transferência de I/O drasticamente inferiores (1-2 MB/s), subutilizando a capacidade de rede e disco do cluster. A implementação de paralelismo via multithreading eleva essa taxa para 20-30 MB/s. Em cenários onde a demanda de I/O ou CPU exceda a capacidade de um único nó, a definição de chaves de sharding adequadas é mandatória para distribuir o esforço de escrita e leitura de forma equitativa entre todos os membros do cluster.