Além da ementa · tema 22 de 30
Spark avançado: DataFrame, tuning e streaming
40 min · 2 vídeos · 14 cards · 1 drill
Por que cai
É o tema em que a sabatina para de perguntar conceito e começa a perguntar diagnóstico. Quem sabe nomear shuffle, skew e broadcast e dizer como confirmaria cada hipótese na Spark UI mostra que já operou job em produção.
Pré-teste · 1 de 2
responda antes de verNa Spark UI, um stage tem 200 tarefas: 199 terminam em 20 segundos e uma leva 25 minutos. Qual a hipótese mais provável?
Vídeo · português · 78 min
Em uma frase
Spark avançado é saber que Catalyst planeja, AQE replaneja e o shuffle é onde o tempo vai embora — e conseguir provar cada hipótese na Spark UI.
DataFrame e Spark SQL: por que a API mudou
Com RDD você descreve como processar; o motor executa o que você mandou. Com DataFrame e Spark SQL você descreve o que quer, e aí o otimizador tem espaço para reescrever. Essa é a razão técnica da mudança, não preferência de sintaxe. Se a banca perguntar quando usar RDD, a resposta honesta é: quase nunca, só em manipulação de baixo nível que a API estruturada não expressa.
Catalyst
Catalyst— Otimizador de consultas do Spark: analisa a expressão, aplica regras de otimização lógica, gera planos físicos candidatos e escolhe um.O que ele faz de mais visível no dia a dia:
- 1Empurra o filtro para perto da leitura (predicate pushdown), inclusive para dentro do Parquet.
- 2Empurra a projeção: se você usa 3 colunas, ele não lê as outras 57.
- 3Poda partições quando o filtro bate na coluna de particionamento.
- 4Reordena joins e escolhe a estratégia de cada um: broadcast, sort-merge ou hash.
- 5Gera código Java especializado para o estágio, em vez de interpretar operador por operador.
O ponto prático: tudo isso depende de o Catalyst enxergar sua lógica. É por isso que UDF em Python custa caro — ela é opaca para o otimizador.
AQE: replanejar durante a execução
O plano do Catalyst nasce de estimativas. AQE usa estatísticas reais dos estágios já concluídos para corrigir o plano no meio do caminho. Três otimizações que valem decorar:
- Junta partições de shuffle pequenas demais, evitando milhares de tarefas minúsculas.
- Divide partições de shuffle desbalanceadas, atacando skew em alguns tipos de join.
- Converte sort-merge join em broadcast join quando descobre que um lado é pequeno de verdade.
Shuffle: o operador caro
Shuffle é redistribuir dados entre executores para que as linhas da mesma chave fiquem juntas. Custa por três motivos somados: serialização e escrita em disco local, transferência pela rede, e a barreira — o próximo estágio só começa quando o anterior termina inteiro.
Provocam shuffle: join por chave, groupBy e agregações, distinct,
orderBy, repartition e window function cuja partição difere da atual.
| repartition | coalesce | |
|---|---|---|
| Faz shuffle | Sim, completo | Não |
| Pode aumentar partições | Sim | Não, só reduz |
| Distribuição resultante | Equilibrada | Pode ficar desigual |
| Uso típico | Rebalancear antes de um join pesado | Reduzir número de arquivos antes de escrever |
Broadcast join
Quando um lado do join é pequeno, o Spark envia a tabela inteira para cada executor e elimina o shuffle do lado grande. Junta transações com a dimensão de bandeira sem embaralhar bilhões de linhas.
Ele escolhe sozinho quando estima que o lado cabe no limite configurado. A estimativa vem das estatísticas da fonte; sem elas, ou com elas desatualizadas, o plano cai em sort-merge join. Quando você vê isso no plano, atualize estatísticas — e só então considere um hint.
Data skew
data skew— Distribuição desigual das chaves, em que uma ou poucas concentram volume desproporcional e caem todas na mesma partição.Como perceber: na Spark UI, dentro de um stage, quase todas as tarefas terminam rápido e uma ou poucas demoram muito, com shuffle read muito maior. Confirme contando linhas por chave.
Como tratar, em ordem de esforço: habilitar o skew join do AQE; usar broadcast, se o outro lado for pequeno; isolar as chaves quentes num caminho separado e unir o resultado; e, por último, salting — acrescentar um sufixo aleatório à chave para espalhá-la em várias partições, replicando o lado menor para cada valor do sal.
Em banco isso é comum: conta agregadora de lojista, cliente institucional com milhões de transações, ou chave nula usada como preenchimento.
Cache e persist
Cache só paga quando o mesmo DataFrame é reutilizado várias vezes numa cadeia cara. Não use quando o dado é lido uma vez só, quando ele não cabe na memória dos executores ou quando ler a origem já é barato. Cache mal colocado rouba memória de execução e provoca spill, deixando o job mais lento exatamente enquanto você tenta acelerá-lo.
UDF em Python
Cada linha atravessa a fronteira JVM-Python com serialização, e o Catalyst não enxerga o que a função faz — logo não empurra filtro, não reordena, não gera código. A ordem de preferência é: função nativa do Spark SQL, depois UDF vetorizada com Pandas (que processa em lote via Arrow), e só então UDF comum.
Structured Streaming, em noções
O stream é tratado como uma tabela que cresce. O motor processa em micro-batches: a cada intervalo de trigger, um pequeno job em lote roda sobre o que chegou desde a última execução. Três conceitos que a banca cobra:
- Trigger: define quando o micro-batch dispara — por intervalo fixo, o mais rápido possível, ou uma única passada sobre o que existe.
- Checkpoint: guarda até onde a fonte foi consumida e o estado das agregações. É o que permite retomar após falha sem perder nem duplicar.
- Modo de saída: append, update ou complete, conforme o que faz sentido para a agregação e para o destino.
Num banco, o caso típico é scoring de fraude sobre autorizações: janela curta, checkpoint obrigatório e destino em tabela Delta para o histórico continuar consultável em lote.
Se quiser outro ângulo
Como cai na sabatina
“Seu job passou de 10 para 40 minutos sem mudar o código. Por onde começa?”
Erros comuns
- Culpar o cluster antes de olhar o plano e a Spark UI. Job que dobra de tempo sem mudança de código quase sempre mudou de dado: volume, distribuição de chave ou quantidade de arquivos.
- Usar repartition onde coalesce bastava. Repartition dispara shuffle completo; coalesce só junta partições existentes e é a escolha certa para reduzir número de arquivos de saída.
- Fazer cache de tudo. Cache ocupa memória do executor, pode empurrar dado para disco e piorar o job; só compensa quando o mesmo DataFrame é reutilizado várias vezes.
- Escrever UDF em Python por comodidade. A serialização entre JVM e Python custa caro e o Catalyst não enxerga dentro da função, então ele não otimiza nada ali.
- Achar que AQE resolve todo skew. AQE trata partições de shuffle desbalanceadas em alguns tipos de join, mas não conserta concentração extrema numa única chave nem skew antes do shuffle.
- Confundir broadcast join com broadcast variable. Broadcast join é a estratégia de enviar a tabela pequena inteira para cada executor; variável de broadcast é outro mecanismo, para dado auxiliar.
- Aumentar shuffle partitions sem olhar o tamanho da partição. Partição pequena demais vira sobrecarga de tarefa; grande demais vira spill em disco.
Flashcards
Card 1 de 14
0 certos · 0 a reverDrill
Diagnostique o job lento
Para cada sintoma, diga a hipótese mais provável e como você a confirmaria. Fale em voz alta antes de conferir, como se estivesse na sabatina.
1.Um stage com 200 tarefas: 199 terminam em 20 segundos, uma leva 25 minutos.
2.O job lê a tabela inteira embora a consulta filtre por uma única data particionada.
3.O job passou de 10 para 40 minutos e ninguém alterou o código.
4.A Spark UI mostra spill de memória para disco no estágio de agregação.
5.Uma tabela de 8 linhas é juntada com bilhões e o plano mostra sort-merge join.
6.O job escreve 40 mil arquivos minúsculos na tabela de saída.
7.Uma transformação simples de string sobre 2 bilhões de linhas domina o tempo do job.
8.O mesmo DataFrame intermediário é usado em cinco saídas diferentes e o job recomputa tudo cinco vezes.
Quiz · 1 de 2
Qual afirmação sobre UDF em Python no Spark é correta?
Explique para um gerente
Explique para um gestor por que redistribuir dados entre máquinas é a parte cara de um processamento distribuído, usando uma analogia sem falar em Spark.
Responda em voz alta antes de seguir. Se travar numa palavra técnica, é sinal de que ainda não entendeu essa parte.
Perguntas de sabatina deste tema
- · Seu job Spark passou de 10 para 40 minutos sem mudança de código. Por onde você começa?
- · Explique o que é shuffle, por que ele é caro e o que você faz para reduzi-lo num pipeline de transações.
Para ir além