Spark e Flink
Nível médioPor 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á.
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.