Spark

Nivel medio

¿Cómo elijo el número correcto de shuffle partitions, memoria del executor y núcleos?

No hay una respuesta universal, pero sí un punto de partida defendible — y, más útil aún, una forma de saber cuándo tus números actuales están mal.

Leer en:
Shuffle partitions Memoria del executor Número de núcleos
Qué dimensionar primero, en orden — no un desglose medido.

Cada guía te da una fórmula. Casi ninguna te dice cómo verificar que la fórmula realmente encajó con tu job. Ambas cosas importan.

Shuffle partitions: parte del tamaño objetivo de partición, no de un número fijo

spark.sql.shuffle.partitions viene por defecto en 200 sin importar el tamaño de los datos, lo cual está mal para casi cualquier carga real — demasiadas particiones para un job pequeño (domina el overhead de programar tareas), muy pocas para uno grande (cada partición hace spill a disco). Un mejor punto de partida: divide el tamaño total de entrada de tu etapa de shuffle entre un objetivo de 128-200MB por partición.

# desde la pestaña Stages de la Spark UI, "Shuffle Read" del stage relevante shuffle_partitions = shuffle_read_bytes / (150 * 1024 * 1024)

Luego verifica: después del cambio, ¿la duración de las tareas dentro de ese stage se volvió más pareja (menos tareas de cola larga) y bajaron los bytes de spill a disco? Si el spill sigue siendo alto, te quedaste corto; si la mayoría de las particiones ahora están casi vacías, te pasaste.

Memoria del executor: dimensiónala según la partición más grande que esperas, no el promedio

Los OOM de executor casi siempre vienen de la partición más grande de un stage, no de la promedio — así que dimensionar la memoria según el tamaño promedio de partición subestima por construcción. Un piso razonable es tamaño_de_la_partición_más_grande * 3-4 (margen para el overhead de deserialización y cualquier agregación en memoria), acotado por lo que el tipo de instancia de tu clúster realmente puede ofrecer por executor.

Número de núcleos: más no es gratis pasado cierto punto

Los núcleos del executor controlan cuántas tareas corren concurrentemente dentro del presupuesto de memoria de un solo executor — llevar esto demasiado alto significa que esas tareas compiten por el mismo heap, lo cual se ve como más tiempo de GC, no más throughput. 4-5 núcleos por executor es un techo razonable para la mayoría de cargas pesadas en JVM; pasar de ahí generalmente cambia presión de GC por un paralelismo que en realidad no ganas.

La métrica que realmente te dice si acertaste: el tiempo de GC como porcentaje del tiempo total de tareas. Por debajo de ~10% es saludable. Subir por encima de 15-20% significa que la presión de memoria está consumiendo tiempo real de cómputo, sin importar cómo hayas configurado tus shuffle partitions o el número de núcleos.

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.