Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.
Lección 2 de 5
Skew y adaptive query execution
Interpreta memory spill y disk spill como síntomas de presión durante sort, aggregate o join, no como una orden automática de comprar más memoria.
- Duración
- 17 min aprox.
- Objetivo
- Interpreta memory spill y disk spill como síntomas de presión durante sort, aggregate o join, no como una orden automática de comprar más memoria.
- Siguiente paso
- Continuar con la siguiente lección
Ver detalles del módulo
Tuning avanzado de Spark
Optimiza con evidencia del plan y las métricas, no con recetas globales ni más cómputo por defecto.
- Diagnosticar skew y spill
- Ajustar joins y particiones
- Evaluar UDF, Pandas UDF y funciones nativas
02ImplementaciónSkew y adaptive query execution
Interpreta memory spill y disk spill como síntomas de presión durante sort, aggregate o join, no como una orden automática de comprar más memoria.
+
Skew y adaptive query execution
Interpreta memory spill y disk spill como síntomas de presión durante sort, aggregate o join, no como una orden automática de comprar más memoria.
Spark derrama datos cuando una operación no puede mantener en memoria sus estructuras intermedias. El spill a disco añade serialización e I/O; un volumen pequeño puede ser normal, pero spill masivo junto con garbage collection, tareas largas o executor lost indica que el plan y la forma de los datos no caben de manera eficiente. Revisa la etapa concreta, el operador y la distribución antes de cambiar la máquina.
Las palancas tienen costes distintos. Reducir el ancho de fila proyectando sólo columnas necesarias, filtrar antes del shuffle o sustituir una UDF por funciones nativas reduce trabajo. Reparticionar puede distribuir mejor la carga, mientras que workers con más memoria ayudan cuando cada partición es legítimamente grande. Aumentar `spark.sql.shuffle.partitions` sin medir puede crear miles de tareas diminutas.
Modelo mental
Piensa en la memoria de ejecución de Spark como una mesa de trabajo compartida, no como un almacén permanente. Un sort, hash aggregate o join construye estructuras temporales por tarea; si la porción activa no cabe, Spark conserva corrección escribiendo parte en disco y leyéndola después. Ese spill es un mecanismo de supervivencia, no necesariamente un fallo. Se vuelve problemático cuando domina el tiempo, coincide con garbage collection, reintentos o executors perdidos, o se concentra en particiones concretas. La pregunta útil no es cuánta memoria tiene todo el clúster, sino qué operador, con qué ancho de fila y distribución, exige cuánta memoria simultánea en cada tarea.
Memoria usada temporalmente por operadores de shuffle, join, sort y aggregate mientras una tarea está activa.
Su presión ocurre por tarea y operador; sumar RAM del clúster no describe si una partición cabe.Estimación de datos intermedios desalojados de estructuras en memoria durante la ejecución.
Señala presión, pero no equivale necesariamente a bytes físicos escritos en disco.Datos intermedios serializados en almacenamiento local para que la operación pueda continuar.
Añade I/O y serialización; un valor alto sostenido suele explicar colas o inestabilidad.from pyspark.sql import functions as F
daily_net = (
events
.where(F.col("event_date") >= F.date_sub(F.current_date(), 30))
.select("event_date", "store_id", "net_amount")
.groupBy("event_date", "store_id")
.agg(F.sum("net_amount").alias("net_amount"))
)
daily_net.explain("formatted")La proyección temprana evita transportar payloads que no participan en la agregación; confirma el cambio en el plan y en los bytes de shuffle.
Puntos clave
- Ubica el spill en una etapa y operador concretos.
- Reduce datos antes del shuffle mediante filtros y proyección de columnas.
- Diferencia presión por partición de falta de memoria global del workload.
Evita
- Aumentar el tamaño de los workers sin eliminar columnas grandes que atraviesan el shuffle.
- Usar `coalesce(1)` para controlar archivos y concentrar toda la escritura en una tarea.
Recuerdo activo