Análise de Desempenho de Clusters MongoDB com Sharding e Replica Sets

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.

Tags: mongodb ycsb Sharding replica-set nosql

Publicado em 9-20 08:30