Spark y Flink
Nivel medio¿Por qué obtengo errores de OutOfMemory y cómo los soluciono de verdad?
"Simplemente aumenta la memoria del executor" arregla el síntoma como la mitad de las veces y desperdicia dinero la otra mitad. Así puedes saber en qué caso estás.
Un error de OOM te dice que se acabó la memoria. No te dice por qué — y la solución es completamente distinta según la causa.
1. Falta de recursos genuina
El caso directo: tu partición/ventana/estado más grande legítimamente no cabe en la memoria que le diste. Esto es real, y la solución sí es más memoria (o una partición más pequeña — ver el artículo de shuffle partitions arriba).
Señal: el uso de heap sube constantemente durante toda la vida de la tarea/job y muere cerca del límite configurado, sin un pico repentino.
2. Una partición sesgada, no un promedio subdimensionado
Si el 95% de tus particiones usan 2GB y una usa 40GB, aumentar la memoria para que todas sobrevivan a ese único valor atípico es caro y a menudo sigue sin ser suficiente. La solución real es atacar el skew en sí — hacer "salting" de una clave de join sesgada, o reparticionar sobre una clave más pareja.
Señal: solo un número pequeño de tareas hace OOM, no la mayoría, y correlacionan con un rango de clave o partición específico.
3. Un broadcast join que salió mal
El umbral de broadcast join de Spark (spark.sql.autoBroadcastJoinThreshold) decide si una tabla más pequeña se copia completa a cada executor. Si esa tabla "pequeña" creció más allá de lo que tus executors pueden retener — o la estimación de tamaño de Spark es incorrecta (común después de un filtro o join anterior que cambia la cardinalidad sin actualizar las estadísticas) — cada executor hace OOM intentando retener un broadcast que ya no es realmente pequeño.
Señal: el OOM ocurre muy temprano en un stage, antes de cualquier trabajo real de shuffle, y desactivar el broadcast join para esa consulta hace que tenga éxito (más lento, pero con éxito).
4. Estado sin límite en un job de streaming (específico de Flink)
En Flink, el OOM a menudo no se trata de un solo batch grande — es estado que nunca debió crecer indefinidamente haciendo exactamente eso: una ventana que nunca cierra por un watermark mal configurado, o estado por clave sin TTL que acumula claves sin fin. Esto no aparecerá en una prueba rápida; aparecerá días o semanas después en producción, con el heap subiendo lentamente y sin recuperarse.
Señal: el tamaño del checkpoint crece constantemente durante días/semanas sin un crecimiento correspondiente en el volumen de entrada.
Antes de aumentar la memoria: revisa primero el tiempo de GC (Spark) o la tendencia de duración de checkpoints (Flink). Ambos son buenas señales de cuál de las cuatro causas estás viendo realmente, y ambos son números que el motor de reglas de opti-pipe ya lee de tu event log o historial de checkpoints, en lugar de que tú tengas que escarbar en la UI de Spark/Flink a mano.
Mira cómo se ve esto en tu propio pipeline.
Sube un event log real de Spark, un run_results.json de dbt, o una exportación de métricas de Flink, y recibe recomendaciones concretas para aprobar — no otra regla general más.