Spark

Nível médio

Como escolho o número certo de shuffle partitions, memória do executor e núcleos?

Não existe uma resposta universal, mas existe um ponto de partida defensável — e, mais útil ainda, uma forma de saber quando seus números atuais estão errados.

Ler em:
Shuffle partitions Memória do executor Número de núcleos
O que dimensionar primeiro, em ordem — não uma divisão medida.

Todo guia te dá uma fórmula. Quase nenhum te diz como verificar se a fórmula realmente serviu para o seu job. As duas coisas importam.

Shuffle partitions: parta do tamanho alvo da partição, não de um número fixo

spark.sql.shuffle.partitions vem 200 por padrão, independente do tamanho dos dados, o que está errado para quase toda carga real — partições demais para um job pequeno (o overhead de agendar tasks domina), partições de menos para um grande (cada partição faz spill em disco). Um ponto de partida melhor: divida o tamanho total de entrada do seu stage de shuffle por um alvo de 128-200MB por partição.

# da aba Stages da Spark UI, "Shuffle Read" do stage relevante shuffle_partitions = shuffle_read_bytes / (150 * 1024 * 1024)

Depois verifique: após a mudança, a duração das tasks dentro daquele stage ficou mais uniforme (menos tasks de cauda longa) e os bytes de spill em disco caíram? Se o spill continua alto, você ficou aquém; se a maioria das partições agora está quase vazia, você passou do ponto.

Memória do executor: dimensione pela maior partição esperada, não pela média

OOMs de executor quase sempre vêm da maior partição de um stage, não da média — então dimensionar memória pela partição média já subestima por construção. Um piso razoável é tamanho_da_maior_partição * 3-4 (margem para o overhead de deserialização e qualquer agregação em memória), limitado pelo que o tipo de instância do seu cluster realmente consegue oferecer por executor.

Número de núcleos: mais não é de graça depois de certo ponto

Os núcleos do executor controlam quantas tasks rodam simultaneamente dentro do orçamento de memória de um único executor — levar isso longe demais faz com que essas tasks disputem o mesmo heap, o que aparece como mais tempo de GC, não mais throughput. 4-5 núcleos por executor é um teto razoável para a maioria das cargas pesadas em JVM; passar disso geralmente troca pressão de GC por um paralelismo que você não ganha de fato.

A métrica que realmente diz se você acertou: tempo de GC como porcentagem do tempo total de tasks. Abaixo de ~10% é saudável. Subir acima de 15-20% significa que a pressão de memória está consumindo tempo real de computação, independentemente de como você configurou suas shuffle partitions ou número de núcleos.

Veja como isso fica no seu próprio pipeline.

Envie um event log real do Spark, um run_results.json do dbt, ou uma exportação de métricas do Flink, e receba recomendações concretas para aprovar — não mais uma regra de bolso.