Spark e Flink

Nível médio

Por que estou recebendo erros de OutOfMemory e como resolvo isso de verdade?

"Simplesmente aumente a memória do executor" resolve o sintoma metade das vezes e desperdiça dinheiro na outra metade. Veja como saber em qual caso você está.

Ler em:
Subprovisionamento genuíno Uma partição enviesada Broadcast join com erro Streaming ilimitado
Ordenado aproximadamente por quão comum é cada causa real de um OOM — não medido.

Um erro de OOM diz que a memória acabou. Não diz por quê — e a correção é completamente diferente dependendo da causa.

1. Subprovisionamento genuíno

O caso direto: sua maior partição/janela/estado legitimamente não cabe na memória que você deu a ela. Isso é real, e a correção realmente é mais memória (ou uma partição menor — veja o artigo de shuffle partitions acima).

Sinal: o uso de heap sobe de forma constante durante toda a vida da task/job e morre perto do limite configurado, sem um pico repentino.

2. Uma partição enviesada, não uma média subdimensionada

Se 95% das suas partições usam 2GB e uma usa 40GB, aumentar a memória para que todas sobrevivam a esse único outlier é caro e muitas vezes ainda insuficiente. A correção real é atacar o skew em si — fazer "salting" de uma chave de join enviesada, ou reparticionar por uma chave mais uniforme.

Sinal: só um número pequeno de tasks dá OOM, não a maioria, e elas se correlacionam com uma faixa de chave ou partição específica.

3. Um broadcast join que deu errado

O limite de broadcast join do Spark (spark.sql.autoBroadcastJoinThreshold) decide se uma tabela menor é copiada inteira para cada executor. Se essa tabela "menor" cresceu além do que seus executors conseguem segurar — ou a estimativa de tamanho do Spark está errada (comum depois de um filtro ou join anterior que muda a cardinalidade sem atualizar estatísticas) — cada executor dá OOM tentando segurar um broadcast que já não é mais pequeno.

Sinal: o OOM acontece bem cedo em um stage, antes de qualquer trabalho real de shuffle, e desativar o broadcast join para aquela query faz ela ter sucesso (mais lento, mas com sucesso).

4. Estado ilimitado em um job de streaming (específico do Flink)

No Flink, o OOM muitas vezes não é sobre um batch grande — é estado que nunca deveria crescer para sempre fazendo exatamente isso: uma janela que nunca fecha por causa de um watermark mal configurado, ou keyed state sem TTL que acumula chaves indefinidamente. Isso não aparece em um teste rápido; aparece dias ou semanas depois em produção, com o heap subindo devagar e não se recuperando.

Sinal: o tamanho do checkpoint cresce de forma constante ao longo de dias/semanas sem crescimento correspondente no volume de entrada.

Antes de aumentar a memória: verifique primeiro o tempo de GC (Spark) ou a tendência de duração de checkpoint (Flink). Ambos são bons indicadores de qual das quatro causas você está realmente vendo, e ambos são números que o motor de regras do opti-pipe já lê do seu event log ou histórico de checkpoints, em vez de você ter que garimpar isso na UI do Spark/Flink na mão.

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.