Spark

Nivel medio

¿Cómo detecto y soluciono el data skew?

El skew es invisible en las métricas agregadas y obvio en las métricas por tarea — solo hay que mirar la vista correcta.

Leer en:
Tarea 1 Tarea 2 Tarea 3 Tarea 4 (la sesgada) Tarea 5
Un stage sesgado, ilustrado — una tarea hace mucho más trabajo que el resto.

La duración promedio de tareas de un job se ve bien. Su duración total no. Esa brecha casi siempre es skew: una partición, o un puñado de ellas, haciendo dramáticamente más trabajo que el resto, con la más lenta dictando el tiempo total del stage.

Cómo verlo de verdad

La duración promedio, e incluso el p95, pueden verse razonables mientras el skew sigue siendo el problema — el skew a menudo afecta una fracción lo suficientemente pequeña de tareas como para que las estadísticas por percentil lo oculten. Lo que no lo oculta: la razón entre la duración máxima de tarea y la duración mediana dentro de un mismo stage. Un stage de shuffle saludable tiene esa razón bien por debajo de 3x. Una razón por encima de 10x es skew, no ruido.

Solución 1: hacer "salting" de la clave sesgada

Para un join o agregación sesgados por un valor de clave específico (un ID de cliente, una región, un flag de estado dominando el volumen), el "salting" distribuye las filas de esa clave entre varias sub-claves sintéticas antes del shuffle, y luego combina los resultados después. Es la solución más quirúrgica cuando sabes cuál clave es el problema, pero requiere un cambio de código en ambos lados de un join.

Solución 2: Adaptive Query Execution (AQE)

Si estás en Spark 3.0+, spark.sql.adaptive.enabled=true (activo por defecto en versiones recientes) incluye manejo automático de skew en joins — Spark detecta una partición sobrecargada en tiempo de ejecución y la divide antes del join, sin necesidad de cambiar código. Esto es lo primero que hay que revisar antes de armar un salting manual: confirmar que AQE realmente está activo y que su umbral de detección de skew no está tan alto que tu caso no lo dispare.

Solución 3: hacer broadcast del lado pequeño en su lugar

Si el skew en realidad es un desbalance de tamaño de join más que un skew dentro de una columna — una tabla es lo suficientemente pequeña para broadcast — cambiar a un broadcast join evita el shuffle (y el skew que viene con él) por completo. Solo es viable si la tabla más pequeña realmente cabe en los límites de memoria del broadcast join; ver el artículo de OOM arriba para lo que pasa cuando no cabe.

Dónde aparece esto en opti-pipe: la varianza de duración por tarea de tu event log de Spark es una de las señales que el motor de reglas lee directamente — cuando la razón máximo/mediana en un stage de shuffle cruza el umbral de skew, se marca como una recomendación específica en lugar de algo que tendrías que notar tú mismo en la Spark UI.

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.