lección 8
Spark UI: diagnosticar por qué tu job tarda 3 horas
El panel de control de un avión: stages, tasks, shuffles, GC y los indicadores que te dicen si algo va mal.
⏱ 55 min
### El panel de instrumentos de tu avión
Imagina que eres piloto. Tu avión tiene cientos de indicadores: altitud, velocidad, temperatura de motores, presión de aceite, nivel de combustible. No necesitas mirarlos todos constantemente — pero cuando algo va mal, esos indicadores te dicen EXACTAMENTE qué está fallando y dónde. La Spark UI es el panel de tu avión de datos.
La Spark UI está disponible en http://localhost:4040 mientras tu job está corriendo (o en el puerto 8080 del master para ver el historial). Es una interfaz web que muestra TODO lo que Spark está haciendo internamente: cuántas stages tiene tu job, cuántas tasks hay en cada stage, cuánto dura cada task, cuántos datos se mueven en cada shuffle, cuánta memoria usa cada executor, y mucho más.
El 90% de los problemas de rendimiento en Spark se diagnostican con la UI. Un job que tarda 3 horas puede pasar a tardar 10 minutos una vez que identificas el cuello de botella. Vamos a aprender a leer los indicadores clave.
### Las pestañas principales
- Jobs: cada acción (count, show, write) genera un job. Muestra el estado y duración de cada uno.
- Stages: cada job se divide en stages (separadas por shuffles). La stage más lenta es tu cuello de botella.
- Storage: DataFrames cacheados/persistidos y cuánta memoria ocupan.
- Environment: configuración de Spark (útil para verificar que tus configs se aplicaron).
- Executors: estado de cada executor — memoria usada, tasks ejecutadas, GC time.
- SQL: planes de ejecución de queries SQL/DataFrame (la vista más útil para optimización).
### Indicador #1: Stage Duration (duración por stage)
Ve a la pestaña Stages y ordena por duración. La stage más lenta domina el tiempo total del job. Si una stage tarda 50 minutos y las demás 2 minutos, tu optimización debe enfocarse en esa stage. Haz clic en ella para ver los detalles.
Dentro de una stage, mira la distribución de tasks: ¿todas las tasks tardan lo mismo o hay una que tarda 10x más? Si una task individual tarda mucho más que las demás, tienes DATA SKEW — una partición tiene muchos más datos que las otras.
### Indicador #2: Shuffle Read/Write
En la vista de stages, las columnas "Shuffle Write" y "Shuffle Read" te dicen cuántos GB se mueven entre stages. Si ves un shuffle de 50 GB en un dataset de 10 GB, algo está mal (¿un join exploding? ¿un cartesian product?). Si el shuffle es proporcionado al tamaño de los datos, es normal pero puedes intentar reducirlo con broadcast joins.
### Indicador #3: GC Time (tiempo de Garbage Collection)
En la pestaña Executors, hay una columna "GC Time". Si un executor pasa más del 10% de su tiempo en GC (Garbage Collection), significa que se está quedando sin memoria. El GC es como un limpiador que intenta liberar memoria descartando objetos que ya no se usan — pero si no hay nada que limpiar y la memoria sigue llena, se convierte en un bucle inútil que ralentiza todo.
Soluciones para GC excesivo: aumentar la memoria por executor (spark.executor.memory), reducir el tamaño de las particiones, o usar menos cache/persist si no es necesario.
### Indicador #4: Spill (desbordamiento a disco)
En los detalles de una stage, las columnas "Spill (Memory)" y "Spill (Disk)" te dicen si los datos intermedios no caben en memoria y se están escribiendo a disco. Spill es muy malo para el rendimiento — significa que el executor está usando el disco como extensión de la RAM, lo que es 100x más lento.
Soluciones: aumentar spark.executor.memory, aumentar spark.memory.fraction (porcentaje de memoria del executor dedicada a datos), o reparticionar para que cada partición sea más pequeña.
### Los 5 problemas más comunes (y cómo los ves en la UI)
- 01.Data Skew: una task tarda 10x más que las demás en la misma stage. Se ve en el gráfico de distribución de duración de tasks.
- 02.Shuffle excesivo: stages con Shuffle Read de decenas de GB. Indica que un broadcast join podría ayudar.
- 03.Pocas particiones: tienes 8 cores pero solo 2 tasks activas. Se ve en la pestaña Executors (cores ociosos).
- 04.Memoria insuficiente: GC Time > 10% o columnas de Spill con valores altos. Necesitas más RAM por executor.
- 05.Small files: Stage 0 (lectura) tarda mucho con miles de tasks muy cortas. Demasiados archivos pequeños en disco.
Consejo de senior: cuando un job de producción se vuelve lento, mi rutina de diagnóstico es: 1) ¿Qué stage es la más lenta? 2) ¿Hay shuffle grande? → broadcast join. 3) ¿Hay una task más lenta? → data skew. 4) ¿Hay spill o GC alto? → más memoria. El 95% de los problemas se resuelven con estos 4 pasos.
### Usando explain() + UI juntos
1from pyspark.sql import SparkSession2from pyspark.sql.functions import col, sum as spark_sum, broadcast34spark = SparkSession.builder.appName("diagnostico").master("spark://spark-master:7077").getOrCreate()5ventas = spark.read.option("header", "true").option("inferSchema", "true").csv("/home/jovyan/data/ventas.csv")6clientes = spark.createDataFrame([(i, f"C{i}") for i in range(1, 50001)], ["id", "nombre"])78# ANTES: ver el plan con explain9print("=== PLAN ANTES DE OPTIMIZAR ===")10resultado_lento = ventas.join(clientes, ventas.cliente_id == clientes.id).groupBy("nombre").agg(spark_sum("cantidad"))11resultado_lento.explain()12# Verás: SortMergeJoin + Exchange (shuffle)1314# Ejecutar y mirar en la UI (http://localhost:4040)15resultado_lento.count() # Mira la duración en Stages1617# DESPUÉS: optimizar con broadcast18print("\n=== PLAN DESPUÉS DE OPTIMIZAR ===")19resultado_rapido = ventas.join(broadcast(clientes), ventas.cliente_id == clientes.id).groupBy("nombre").agg(spark_sum("cantidad"))20resultado_rapido.explain()21# Verás: BroadcastHashJoin, sin Exchange en ventas2223resultado_rapido.count() # Compara la duración en Stages24# Spoiler: la stage del join pasará de minutos a segundos
El workflow de optimización: explain() para ver el plan, UI para medir el impacto real
La Spark UI del driver (puerto 4040) solo está disponible MIENTRAS el job corre. Para ver el historial de jobs pasados, activa el Spark History Server o mira la UI del master (puerto 8080). En producción siempre tendrás un History Server configurado.
1# Para que la UI siga ahí cuando el job termine, activa el event log.2# Sin esto, el puerto 4040 desaparece con todos sus números en cuanto3# el script acaba — y para un job de 30 segundos, eso es no verlos.4spark = (5 SparkSession.builder6 .appName("diagnostico")7 .master("spark://spark-master:7077")8 .config("spark.eventLog.enabled", "true")9 .config("spark.eventLog.dir", "/home/jovyan/output/spark-events")10 .getOrCreate()11)12# Spark deja en esa carpeta un fichero JSON con TODO lo que la UI mostraba.13# El History Server lee ese fichero; tú también puedes abrirlo a mano.
Dos líneas para que la Spark UI sobreviva al final del job
### Configuraciones clave para optimización
1# Las configs más importantes para rendimiento2spark.conf.set("spark.sql.shuffle.partitions", "100") # default 200, ajustar al dataset3spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "50m") # broadcast hasta 50 MB4spark.conf.set("spark.sql.adaptive.enabled", "true") # Adaptive Query Execution (Spark 3+)5spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true") # reduce particiones vacías67# Adaptive Query Execution (AQE) — la joya de Spark 3.08# AQE ajusta el plan de ejecución EN TIEMPO REAL basándose en estadísticas reales:9# - Reduce particiones vacías automáticamente10# - Detecta y maneja data skew11# - Elige broadcast join si los datos son más pequeños de lo estimado12# Siempre actívalo. No hay razón para no hacerlo.13print("AQE habilitado:", spark.conf.get("spark.sql.adaptive.enabled"))
Configuraciones que deberías ajustar en cada job de producción
Consejo de senior: Adaptive Query Execution (AQE) fue el cambio más impactante de Spark 3.0. Antes necesitabas semanas de tuning manual. Ahora Spark se auto-optimiza en muchos casos. Siempre actívalo. Es como poner tu coche en modo "adaptativo" — ajusta la suspensión según el terreno automáticamente.
## ejercicios
Diagnosticar un job lento
Este pipeline tiene un problema de rendimiento. Ejecútalo, mira la Spark UI, identifica el cuello de botella y corrígelo. El pipeline hace un join de ventas con clientes (tabla pequeña) usando un shuffle join innecesario.
Activar AQE y comparar rendimiento
Ejecuta el mismo job con AQE desactivado y activado. Observa cómo AQE reduce automáticamente las particiones vacías y optimiza el plan.
Checklist de optimización para producción
Escribe una función que reciba las métricas de un job (duración de stages, shuffle size, GC time, spill) y devuelva una lista de recomendaciones de optimización ordenadas por impacto.
Regístrate para guardar tu progreso.
## comentarios
Reporta erratas, ayuda a otros o comparte tu opinión. Sé constructivo.
Inicia sesión para comentar y responder.
cargando comentarios...