Spark

Junior

¿Por qué mi trabajo de Spark de repente se volvió más lento sin cambios en el código?

Nueve de cada diez veces no es el código. Son los datos, el clúster o una configuración que dejó de coincidir con alguno de los dos, silenciosamente.

Leer en:
Crecieron los datos Join desbalanceado Cluster distinto al probado
Ordenado aproximadamente por qué tan seguido es la causa real — no medido en un dataset real.

No tocaste el job. El DAG es idéntico. El último deploy fue hace tres semanas. Y aun así, la corrida de anoche tardó 40 minutos en lugar de 12. Este es uno de los tickets más comunes en data engineering, y casi siempre se reduce a una de tres causas.

1. Los datos crecieron, pero la configuración no

El número de shuffle partitions, la memoria del executor y el umbral de broadcast join suelen fijarse una sola vez, durante el desarrollo inicial, según el volumen de datos de ese momento. Seis meses después, la misma tabla tiene 4 veces más filas y el mismo spark.sql.shuffle.partitions=200, que ahora produce particiones demasiado grandes para caber cómodamente en memoria, forzando spills a disco en cada etapa de shuffle. El código no cambió; lo que cambió fue la suposición que la configuración codificaba.

Cómo verificarlo: compara el conteo de filas de entrada y los bytes de lectura/escritura de shuffle entre una corrida lenta reciente y una antigua rápida, en la pestaña Stages de la Spark UI. Un salto grande ahí, con una desaceleración proporcional, apunta directamente a esta causa.

2. Un join que antes estaba balanceado ya no lo está

El data skew no se anuncia — simplemente significa que un puñado de particiones terminan haciendo 10-100 veces más trabajo que el resto, así que el tiempo total lo dicta la tarea más lenta, no el promedio. Una distribución de claves que era razonablemente uniforme al lanzar puede derivar conforme cambian los patrones de uso (un cliente, una región, un tipo de evento empieza a dominar el volumen).

Cómo verificarlo: en la vista de detalle del stage en la Spark UI, ordena las tareas por duración. Un stage donde la mediana es 4 segundos y el máximo es 6 minutos está sesgado, sin ninguna duda.

3. El clúster ya no es el clúster con el que probaste

Los clústeres con autoescalado, los pools de instancias spot y los clústeres compartidos multi-tenant pueden asignarle silenciosamente a tu job menos executors, o más lentos, especialmente bajo contención con otros jobs. Esto se ve como un tiempo total más lento con un tiempo de CPU por tarea idéntico — el job hace la misma cantidad de trabajo, solo espera más por los recursos para hacerlo.

Cómo verificarlo: compara el número de executors realmente asignados (no solicitados) entre corridas, y mira el retraso del scheduler de tareas, no solo la duración de las tareas.

La forma más rápida de distinguir entre estas causas: necesitas dos cosas lado a lado — tu configuración actual y métricas reales de la corrida que realmente ocurrió. Esa es la premisa completa del motor de reglas de opti-pipe: lee el tiempo real de las tareas, el volumen de shuffle y el uso de memoria de tu event log, y lo compara con tus shuffle partitions, memoria de executor y número de instancias configurados — en lugar de que tú mires la Spark UI y adivines cuál de las tres causas es.

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.