Spark

Nível médio

Como detecto e corrijo data skew?

O skew é invisível nas métricas agregadas e óbvio nas métricas por task — basta olhar a visão certa.

Ler em:
Task 1 Task 2 Task 3 Task 4 (a enviesada) Task 5
Um stage enviesado, ilustrado — uma task fazendo muito mais trabalho que as outras.

A duração média das tasks de um job parece boa. A duração total não. Essa diferença quase sempre é skew: uma partição, ou um punhado delas, fazendo dramaticamente mais trabalho que as demais, com a mais lenta ditando o tempo total do stage.

Como realmente enxergar isso

A duração média, e até o p95, podem parecer razoáveis enquanto o skew continua sendo o problema — o skew muitas vezes afeta uma fração pequena o suficiente de tasks para que estatísticas de percentil o escondam. O que não esconde: a razão entre a duração máxima de task e a duração mediana dentro de um mesmo stage. Um stage de shuffle saudável tem essa razão bem abaixo de 3x. Uma razão acima de 10x é skew, não ruído.

Correção 1: fazer "salting" da chave enviesada

Para um join ou agregação enviesados por um valor de chave específico (um ID de cliente, uma região, uma flag de status dominando o volume), o salting distribui as linhas daquela chave por várias sub-chaves sintéticas antes do shuffle, e depois combina os resultados. É a correção mais cirúrgica quando você sabe qual chave é o problema, mas exige mudança de código nos dois lados de um join.

Correção 2: Adaptive Query Execution (AQE)

Se você está no Spark 3.0+, spark.sql.adaptive.enabled=true (ativo por padrão nas versões recentes) inclui tratamento automático de skew em joins — o Spark detecta uma partição sobrecarregada em tempo de execução e a divide antes do join, sem mudança de código. Isso é a primeira coisa a verificar antes de montar um salting manual: confirmar que o AQE está realmente ativo e que o limite de detecção de skew dele não está tão alto que seu caso não o dispara.

Correção 3: fazer broadcast do lado pequeno em vez disso

Se o skew na verdade é um desequilíbrio de tamanho de join em vez de um skew dentro de uma coluna — uma tabela é pequena o suficiente para broadcast — trocar para um broadcast join evita o shuffle (e o skew junto com ele) por completo. Só é viável se a tabela menor realmente couber nos limites de memória do broadcast join; veja o artigo de OOM acima para o que acontece quando não cabe.

Onde isso aparece no opti-pipe: a variância de duração por task do seu event log do Spark é um dos sinais que o motor de regras lê diretamente — quando a razão máximo/mediana em um stage de shuffle cruza o limite de skew, isso é sinalizado como uma recomendação específica em vez de algo que você teria que notar sozinho na Spark UI.

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.