Saltar al contenido
Lakehouse LabLakehouse LabPreparación Databricks Data Engineer
Módulo 23 · Professional

Tuning Spark

Contenido abierto

Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.

Professional

Tuning avanzado de Spark

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

Lectura pública
Al terminar podrás
  • Diagnosticar skew y spill
  • Ajustar joins y particiones
  • Evaluar UDF, Pandas UDF y funciones nativas
Ver fuentes y revisión

Metadatos editoriales

Última revisión
21 jul 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
01
Modelo mental

Plan físico y métricas de stages

Aprende a distinguir skew real de una etapa simplemente costosa usando la distribución de tiempos, bytes y registros por tarea.

Objetivo
Aprende a distinguir skew real de una etapa simplemente costosa usando la distribución de tiempos, bytes y registros por tarea.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

En Spark, una etapa termina cuando acaba su tarea más lenta. Por eso una media razonable puede ocultar una cola extrema: si la mediana dura 18 segundos y una tarea tarda 11 minutos, añadir workers no elimina la clave caliente que concentra datos. La evidencia útil está en la pestaña Stages de Spark UI: duración máxima frente a mediana, registros de entrada, shuffle read y tamaño de cada tarea.

Antes de aplicar salting, confirma la causa. Un join con una clave nula dominante, un cliente desproporcionado o una partición temporal demasiado amplia producen remedios distintos. AQE puede dividir particiones sesgadas en joins compatibles; si la distribución forma parte del negocio, conviene además aislar claves calientes o rediseñar la agregación y medir el efecto con el mismo conjunto de datos.

Modelo mental

Imagina una etapa de Spark como una carrera por equipos cuyo tiempo oficial lo determina el último corredor. Cada partición produce una tarea y todas deben terminar antes de avanzar; por eso el promedio oculta el dato decisivo. El skew no significa simplemente que el conjunto sea grande, sino que el reparto de trabajo entre tareas es muy desigual. Una clave frecuente, un valor nulo dominante o un rango temporal desproporcionado concentra registros y bytes en pocas particiones. El diagnóstico correcto enlaza tres niveles: distribución del dato de negocio, particionado físico tras el exchange y cola de duraciones observada en Spark UI. Escalar compute mejora capacidad general, pero no redistribuye una clave caliente por sí solo.

Skew de datos

Distribución en la que unas pocas claves o rangos concentran una fracción desproporcionada de registros y trabajo.

Convierte unas pocas tareas en el camino crítico aunque el clúster tenga capacidad ociosa.
Exchange

Frontera del plan físico que redistribuye datos entre executors, normalmente para joins, agregaciones o ventanas.

Es el punto donde la distribución lógica de claves se materializa como particiones y puede revelar skew.
Straggler

Tarea mucho más lenta que sus pares dentro de la misma etapa.

La etapa no finaliza hasta que termina el straggler; por ello el máximo pesa más que la media.
PySparkMedir la distribución de una clave antes del join
from pyspark.sql import functions as F

key_profile = (
    orders.groupBy("customer_id")
    .count()
    .orderBy(F.desc("count"))
)

key_profile.show(20, truncate=False)
orders.where(F.col("customer_id").isNull()).count()

La tabla de frecuencias no sustituye Spark UI, pero permite conectar una tarea extrema con una clave de negocio concreta.

Puntos clave

  • Compara percentiles y máximos por tarea, no sólo la duración media de la etapa.
  • Relaciona la tarea lenta con sus bytes de shuffle y número de registros.
  • Corrige la distribución de datos antes de aumentar capacidad de forma permanente.

Evita

  • Confundir muchas tareas pequeñas con skew: en ese caso el problema puede ser sobreparticionado y overhead de planificación.
  • Aplicar salting a todas las claves y encarecer el join aunque sólo una fracción mínima esté sesgada.

Recuerdo activo

Una etapa tiene 4.000 tareas; 3.995 duran menos de 25 segundos y cinco superan 12 minutos con diez veces más shuffle read. ¿Cuál es la primera hipótesis?

Borrador privado · solo en este navegador
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.

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.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

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.

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

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

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

Borrador privado · solo en este navegador
03
Operación

Broadcast y sort merge join

Comprende qué puede reoptimizar Adaptive Query Execution en tiempo de ejecución y qué decisiones siguen dependiendo del diseño del ingeniero.

Objetivo
Comprende qué puede reoptimizar Adaptive Query Execution en tiempo de ejecución y qué decisiones siguen dependiendo del diseño del ingeniero.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

AQE usa estadísticas disponibles después de exchanges para cambiar una estrategia sort-merge a broadcast hash, combinar particiones post-shuffle demasiado pequeñas, dividir particiones sesgadas y propagar relaciones vacías. En Databricks está habilitado por defecto para consultas batch compatibles con exchanges o subconsultas. El plan adaptativo final puede diferir del plan inicial mostrado antes de ejecutar.

AQE no reordena dinámicamente todos los joins ni arregla un modelo de datos deficiente. Una relación que parece pequeña en catálogo puede superar el límite real, y determinados tipos de join no admiten broadcast en uno de sus lados. Usa `explain('formatted')`, el plan final de Spark UI y métricas de ejecución para demostrar qué regla se aplicó; evita copiar configuraciones antiguas que anulen los defaults optimizados.

Modelo mental

Catalyst prepara una ruta con estimaciones; AQE actúa como un navegador que recalcula cuando ya conoce el tráfico real después de ciertos cruces. Esos cruces son query stages separados por exchanges o subconsultas. Al materializarse una etapa, Spark obtiene tamaños y distribuciones más fiables que las estadísticas previas y puede cambiar algunas decisiones físicas sin alterar la consulta lógica. AQE no es un optimizador omnisciente: no corrige semántica, no rediseña el modelo ni reordena libremente toda cadena de joins. Su valor está en adaptar particiones post-shuffle, tratar skew, propagar relaciones vacías y, cuando procede, sustituir un sort-merge por broadcast con evidencia runtime.

Query stage

Fragmento del plan adaptativo delimitado por exchanges cuya salida puede materializar estadísticas runtime.

AQE toma nuevas decisiones con información precisa al terminar cada etapa materializable.
Post-shuffle coalescing

Combinación dinámica de particiones pequeñas producidas por un shuffle.

Reduce overhead de tareas diminutas sin imponer un número estático adecuado para todos los volúmenes.
Plan final adaptativo

Plan físico efectivo después de aplicar o descartar reglas AQE durante la ejecución.

Es la evidencia para saber qué estrategia se usó; el plan inicial no basta.
PySparkInspeccionar configuración y plan adaptativo
print(spark.conf.get("spark.sql.adaptive.enabled"))

result = (
    facts.join(dimensions, "product_id")
    .groupBy("category")
    .count()
)

result.explain("formatted")
result.count()  # materializa el plan para revisarlo en Spark UI

La acción materializa la consulta; revisa después el plan final y no deduzcas la estrategia sólo del plan inicial.

Puntos clave

  • AQE decide con estadísticas posteriores al shuffle, más precisas que muchas estimaciones previas.
  • Puede coalescer particiones y tratar skew sin alterar el resultado lógico.
  • No sustituye el orden lógico de joins ni una buena reducción temprana de datos.

Evita

  • Suponer que AQE reordena automáticamente una cadena de joins mal diseñada.
  • Desactivar AQE para reproducir una configuración heredada sin comparar resultados y métricas.

Recuerdo activo

¿Qué ventaja tiene un broadcast planeado estáticamente frente a uno elegido tarde por AQE?

Borrador privado · solo en este navegador
04
Diagnóstico

Shuffle partitions y file sizing

Selecciona broadcast, sort-merge u otra estrategia según tamaño, tipo de join, estadísticas y riesgo para el driver y los executors.

Objetivo
Selecciona broadcast, sort-merge u otra estrategia según tamaño, tipo de join, estadísticas y riesgo para el driver y los executors.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

Un broadcast hash join replica la relación pequeña y evita repartir la grande por la clave. Es excelente para una dimensión realmente pequeña, pero peligroso si las estadísticas están obsoletas o si el lado difundido crece: la transferencia y la tabla hash consumen memoria en cada executor. Los hints expresan una preferencia al optimizador, no corrigen una semántica de join incompatible.

Para joins grandes, el sort-merge distribuye ambos lados y paga shuffle y ordenación. Reduce primero las filas y columnas, conserva estadísticas y observa si hay claves calientes. En SQL y DataFrames, el criterio no es 'broadcast siempre es más rápido', sino coste total, compatibilidad del tipo de join y evidencia estable en ejecuciones representativas.

Modelo mental

Una estrategia de join decide dónde se encuentran las filas, cuánto dato viaja y qué memoria se replica. Broadcast lleva una relación pequeña a cada executor y evita redistribuir la grande; sort-merge redistribuye ambos lados por clave y los ordena; shuffle hash también reparte y construye mapas por partición. No existe una estrategia universalmente rápida. La elección depende del tamaño después de filtros y proyecciones, el tipo de join, la distribución de claves, la calidad de estadísticas y los límites de memoria. El modelo mental correcto compara coste total y riesgo: red, ordenación, memoria repetida, posible skew y estabilidad cuando el lado supuestamente pequeño crece.

Build side

Lado del join usado para construir la estructura hash, local o difundida.

Determina memoria, compatibilidad con outer joins y qué conjunto debe permanecer acotado.
Broadcast exchange

Operación que recopila y distribuye una relación a los executors que ejecutan el join.

Evita un shuffle grande, pero replica bytes y puede fallar si la relación crece.
Sort-merge join

Estrategia que reparte ambos lados por clave, los ordena y fusiona secuencialmente.

Escala a relaciones grandes a cambio de red, ordenación, memoria temporal y posible spill.
PySparkBroadcast explícito de una dimensión acotada
from pyspark.sql import functions as F

active_products = (
    spark.table("prod.ref.products")
    .where("is_active = true")
    .select("product_id", "category")
)

enriched = orders.join(F.broadcast(active_products), "product_id", "left")
enriched.explain("formatted")

Documenta el tamaño máximo esperado de `active_products`; un snapshot pequeño hoy no garantiza que siga siendo difundible.

Puntos clave

  • Difunde sólo relaciones acotadas cuyo tamaño conoces en producción.
  • Un hint no cambia qué lado puede difundirse en cada tipo de join.
  • Actualiza estadísticas y compara el plan físico, no sólo el código fuente.

Evita

  • Forzar broadcast a partir de un `count()` de muestra que no representa el máximo diario.
  • Difundir una tabla ancha cuando bastaba proyectar dos columnas de la dimensión.

Recuerdo activo

Una dimensión tiene 20 millones de filas pero sólo se necesitan dos columnas y los registros activos son el 1 %. ¿Qué harías antes del join?

Borrador privado · solo en este navegador
05
Decisión de diseño

UDF, Pandas UDF y serialización

Evita barreras entre Python y el motor usando expresiones nativas y reserva UDFs para lógica que la plataforma no puede expresar.

Objetivo
Evita barreras entre Python y el motor usando expresiones nativas y reserva UDFs para lógica que la plataforma no puede expresar.
Duración estimada
17 min aprox.
Dificultad
Professional
Prerrequisitos
m12
Reportar un error en esta lección

Las funciones SQL y PySpark nativas permanecen visibles para Catalyst y pueden beneficiarse de Photon, generación de código, pushdown y simplificación de expresiones. Una UDF escalar de Python serializa datos entre JVM y Python, oculta parte de la lógica al optimizador y puede provocar fallback de Photon. Antes de crearla, busca funciones integradas para arrays, mapas, strings, fechas y tipos complejos.

Cuando la lógica no existe, una pandas UDF o APIs basadas en Arrow pueden procesar lotes y reducir el coste por fila, pero siguen requiriendo medición, tipos explícitos y pruebas de nulos. La optimización correcta incluye mantenibilidad: una expresión nativa legible suele ser más fácil de gobernar y portar que una UDF opaca.

Modelo mental

Catalyst sólo puede optimizar lo que entiende. Una expresión nativa forma parte del árbol lógico: el motor conoce tipos, nulabilidad y operadores, puede plegar constantes, empujar filtros y ejecutar con Photon cuando hay soporte. Una UDF de Python se parece a una caja negra situada al otro lado de una frontera de proceso; Spark debe serializar columnas, transferir lotes o filas y aceptar que no puede razonar sobre la lógica interna. Esto no hace ilegítimas las UDF, pero cambia la carga de prueba. Primero se buscan funciones SQL, funciones de orden superior y operaciones de tipos complejos; sólo la necesidad funcional justifica perder visibilidad y añadir contrato explícito.

Expresión Catalyst

Nodo tipado del plan lógico que representa una operación conocida por el optimizador.

Permite simplificación, pushdown y elección de operadores nativos o Photon.
Frontera Python

Transferencia y serialización entre el proceso que ejecuta Spark y un worker Python.

Añade coste por lote o fila y oculta la semántica interna al optimizador.
Función de orden superior

Función nativa que transforma o filtra elementos de arrays y mapas mediante expresiones lambda SQL.

Resuelve lógica compleja manteniéndola visible y optimizable, a menudo evitando UDFs.
PySparkNormalización nativa sin UDF
from pyspark.sql import functions as F

normalized = customers.withColumn(
    "email_domain",
    F.lower(F.element_at(F.split(F.trim("email"), "@"), -1)),
).withColumn(
    "is_company_email",
    ~F.col("email_domain").isin("gmail.com", "outlook.com", "yahoo.com"),
)

La expresión queda disponible para el plan; prueba explícitamente emails nulos, sin arroba y con espacios.

Puntos clave

  • Prefiere funciones nativas porque el optimizador conserva visibilidad de la expresión.
  • Comprueba en Query Profile o Spark UI si una UDF provoca fallback de Photon.
  • Si necesitas Python, vectoriza por lotes y define contratos de tipos y nulos.

Evita

  • Crear una UDF para operaciones ya cubiertas por `when`, `transform`, `regexp_extract` o funciones de fecha.
  • Sustituir una UDF escalar por pandas UDF sin medir serialización, tamaño de lote y presión de memoria.

Recuerdo activo

¿Qué señal confirma que una UDF perjudica Photon?

Borrador privado · solo en este navegador
5 lecciones pendientes

Vista de lectura · sin ejecución

Learn Databricks

Learn Databricks · commit 08c378c

README.md