Spark

Nível médio

Quando o broadcast join do Spark realmente ajuda, e quando ele se volta contra você?

Um broadcast join é uma aposta de que um lado do join é pequeno o suficiente para copiar para todo lugar barato. Quando a aposta acerta, é o join mais rápido que o Spark tem. Quando erra, é a forma mais rápida de derrubar um driver por OOM.

Ler em:
Broadcast correto Estimativa de tamanho antiga Filtro não aplicado
Aproximadamente com que frequência cada cenário explica uma surpresa de broadcast join - ilustrativo, não medido.

spark.sql.autoBroadcastJoinThreshold (10MB por padrão) diz ao planejador: se o tamanho estimado de uma tabela está abaixo disso, pule o shuffle e mande uma cópia completa para cada executor.

O que o threshold realmente controla

Abaixo do threshold, o planejador troca um shuffle join (os dois lados reparticionados e embaralhados pela rede) por um broadcast join (o lado pequeno coletado no driver, depois enviado a cada executor como tabela hash). Sem shuffle para o lado grande significa sem stage de shuffle read, sem o skew associado ao shuffle, e normalmente um join visivelmente mais rápido - exatamente por isso o planejador faz isso automaticamente sempre que estima que uma tabela é pequena o suficiente.

O modo de falha: um broadcast que deveria ter ficado como shuffle

A estimativa que o planejador usa vem das estatísticas da tabela, que ficam desatualizadas - uma tabela que tinha 8MB quando as estatísticas foram calculadas pela última vez mas agora tem 800MB ainda vai ser enviada por broadcast se ninguém atualizou essas estatísticas. Resultado: cada executor tenta manter uma tabela hash de 800MB em memória ao mesmo tempo, e muitas vezes é o driver (que coleta os dados do broadcast primeiro) que dá OOM antes de qualquer executor.

Como identificar isso pelo event log

Procure um stage com pouquíssimas tarefas (uma etapa de coleta de broadcast basicamente roda em uma única tarefa) mostrando um uso de memória de pico incomumente alto, seguido imediatamente por uma falha em nível de stage ou um OOM do driver nos logs da aplicação, em vez de uma falha em nível de tarefa. Um OOM de shuffle join, por outro lado, aparece como muitas tarefas falhando num stage de tamanho normal - uma assinatura completamente diferente.

Se não tiver certeza se um broadcast está acontecendo: confira o plano físico procurando BroadcastHashJoin versus SortMergeJoin - é o único lugar onde o Spark diz diretamente qual estratégia escolheu, em vez de deixar você inferir pelo tempo.

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 que exigem sua aprovação - não mais uma regra geral.