Pular para o conteúdo

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 ver

Na Spark UI, um stage tem 200 tarefas: 199 terminam em 20 segundos e uma leva 25 minutos. Qual a hipótese mais provável?

Confiança:

Vídeo · português · 78 min

É o vídeo mais próximo de uma sabatina de tuning em português: diagnóstico guiado, não lista de dicas. Com pouco tempo, vá direto aos blocos de shuffle, skew e join, que são os que a banca cobra.

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

CatalystOtimizador 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:

  1. 1Empurra o filtro para perto da leitura (predicate pushdown), inclusive para dentro do Parquet.
  2. 2Empurra a projeção: se você usa 3 colunas, ele não lê as outras 57.
  3. 3Poda partições quando o filtro bate na coluna de particionamento.
  4. 4Reordena joins e escolhe a estratégia de cada um: broadcast, sort-merge ou hash.
  5. 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.

repartitioncoalesce
Faz shuffleSim, completoNão
Pode aumentar partiçõesSimNão, só reduz
Distribuição resultanteEquilibradaPode ficar desigual
Uso típicoRebalancear antes de um join pesadoReduzir 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 skewDistribuiçã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

Cinco minutos para fechar o conceito de replanejamento em tempo de execução. Guarde as três otimizações que ele cita: elas são a resposta pronta para a pergunta sobre AQE.

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 rever

Drill

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. 1.Um stage com 200 tarefas: 199 terminam em 20 segundos, uma leva 25 minutos.

  2. 2.O job lê a tabela inteira embora a consulta filtre por uma única data particionada.

  3. 3.O job passou de 10 para 40 minutos e ninguém alterou o código.

  4. 4.A Spark UI mostra spill de memória para disco no estágio de agregação.

  5. 5.Uma tabela de 8 linhas é juntada com bilhões e o plano mostra sort-merge join.

  6. 6.O job escreve 40 mil arquivos minúsculos na tabela de saída.

  7. 7.Uma transformação simples de string sobre 2 bilhões de linhas domina o tempo do job.

  8. 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.
Responder no simulado

Para ir além