Spark

Nivel medio

¿Cómo configuro dynamic allocation de Spark desde cero, no solo la ajusto después?

Dynamic allocation suele discutirse como una perilla de ajuste, pero tiene requisitos previos duros - sin ellos, silenciosamente no hace nada, lo cual se ve idéntico a un job mal ajustado.

Leer en:
Shuffle service externo Límites min/max de executors Timeout de inactividad
Aproximadamente el orden en que estas configuraciones necesitan atención - el shuffle service no es opcional, el resto es ajuste sobre esa base.

spark.dynamicAllocation.enabled=true por sí solo no hace nada en la mayoría de los cluster managers. Dynamic allocation primero necesita un shuffle service externo configurado y corriendo - sin él, los executors no se pueden quitar de forma segura a mitad del job, así que Spark simplemente... no los quita.

El requisito previo que nadie menciona primero

Quitar un executor a mitad del job solo es seguro si los datos de shuffle que contiene siguen disponibles para otros executors después de que se va - eso es lo que hace el shuffle service externo, corriendo independientemente del ciclo de vida de cualquier executor individual. Sin él habilitado, Spark acepta dynamicAllocation.enabled=true sin error y luego simplemente nunca reduce la escala, porque hacerlo arriesgaría perder datos de shuffle. En YARN y Kubernetes esto es un servicio separado que hay que desplegar y al que hay que apuntar desde la configuración de Spark - no aparece automáticamente solo porque dynamic allocation está encendido.

# spark-defaults.conf spark.shuffle.service.enabled true spark.dynamicAllocation.enabled true

Establecer límites reales, no solo un techo

Dejar minExecutors en su valor por defecto de 0 significa que un job puede reducirse por completo a cero executors durante un momento de calma y luego pagar la latencia completa de asignación del cluster manager para volver a escalar - está bien para un job por lotes con margen, costoso para cualquier cosa sensible a la latencia. maxExecutors es la perilla más familiar (es el techo de costo), pero es solo la mitad del panorama sin un piso correspondiente:

# spark-defaults.conf spark.dynamicAllocation.minExecutors 2 spark.dynamicAllocation.maxExecutors 40 spark.dynamicAllocation.initialExecutors 2

Timeout de inactividad, y dónde todavía hace falta criterio

spark.dynamicAllocation.executorIdleTimeout (por defecto 60s) controla cuánto tiempo un executor permanece inactivo antes de liberarse - demasiado corto y un job con etapas irregulares y en ráfagas sacude executors hacia arriba y abajo constantemente, agregando latencia de asignación en cada ráfaga; demasiado largo y está pagando por capacidad inactiva entre ráfagas. No hay una fórmula para el valor correcto - depende de cuán irregular sea realmente la carga de trabajo, que es exactamente el tipo de patrón que vale la pena mirar en la línea de tiempo de un run real antes de elegir un número.

Dónde encaja opti-pipe: el motor de reglas lee el event log de un job completado y señala indicios de sub- o sobre-aprovisionamiento contra los límites que realmente están configurados - no establece minExecutors/maxExecutors por usted, ya que eso es un compromiso entre costo y latencia que solo usted puede hacer.

Vea cómo se ve esto en su propio pipeline.

Suba un event log real de Spark, un run_results.json de dbt o una exportación de métricas de Flink y obtenga recomendaciones concretas, que debe aprobar antes de aplicar - no otra regla general.