Saltar al contenido

Tuning Spark

Menú

Puedes leer sin crear un espacio. Créalo solo cuando quieras guardar.

Guardar 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.

17 min aprox.

Detalles

Tuning avanzado de Spark

Optimiza con evidencia del plan y las métricas, no con recetas globales ni más cómputo por defecto.

Reto observable

Optimiza un workload sesgado con una hipótesis medible y presenta plan, shuffle, spill, duración y corrección antes y después.

Al terminar podrás
  • Diagnosticar skew y spill
  • Ajustar joins y particiones
  • Evaluar UDF, Pandas UDF y funciones nativas
Prerrequisitos
m12
Última revisión
25 ago 2026
Nivel
Professional
Ruta relacionada
performance
Dominios blueprint
Cost & Performance Optimisation · Spark UI
Estado
Revisión editorial interna
Fuentes principales
Adaptive query execution · Diagnose cost and performance issues using the Spark UI
Reportar un error
02
Implementación

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.

PySparkReducir el ancho antes de agregar
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.

¿Por qué un spill elevado no demuestra por sí solo que falte memoria en todo el clúster?

Profundiza

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.

Execution memory

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.
Memory spill

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.
Disk spill

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.
Resumen

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.

Vista de lectura · sin ejecución

Learn Databricks

Learn Databricks · commit 08c378c

README.md

Módulo 23

Contenido del módulo