Reducción de escala eficiente y gestor de shuffle remoto

Se aplica a:✅ Ingeniería de datos de tejido y ciencia de datos

La reducción de escala eficiente es una característica de Spark de Microsoft Fabric que desacopla los datos de barajado de Spark del ciclo de vida del ejecutor. En lugar de fijar la salida del shuffle en los discos locales del ejecutor, Fabric Spark enruta los datos del shuffle a Azure Blob Storage (o los migra allí bajo demanda) y permite que la Ejecución adaptable de consultas (AQE) determine cómo se realiza la propia escritura. El resultado es una reducción de escala del clúster más rápida, un menor coste de computación y trabajos más resilientes, sin cambios en tus consultas, notebooks ni pipelines.

Overview

El escalado a la baja eficiente se basa en cuatro capacidades que cooperan entre sí:

Capacidad Qué hace
Gestor remoto de Shuffle (RSM) Escribe y lee los datos de shuffle en Azure Blob Storage en lugar de en los discos locales del ejecutor.
Migración aleatoria Mueve los bloques de mezcla fuera de un ejecutor antes de retirarlo, en lugar de descartarlos.
Capa de decisión Enrutamiento en tiempo de ejecución por fase que mantiene pequeñas ordenaciones aleatorias locales y descarga ordenaciones aleatorias grandes al almacenamiento remoto.
Escritura aleatoria de AQE Permite que Adaptive Query Execution participe en la fase de escritura de shuffle, para que el particionado sea correcto desde el primer momento.

Prerequisites

  • El motor de ejecución nativo (NEE) debe estar habilitado.
  • Escalado automático habilitado (recomendado). La reducción de escala eficiente también funciona sin escalado automático mediante las siguientes configuraciones de Spark.
  • Runtime 1.3 (Apache Spark 3.5) o posterior.

Cómo funciona

Cuando Spark procesa una consulta, a menudo redistribuye datos entre fases: un shuffle. Normalmente, los datos de mezcla se almacenan en el disco local de cada ejecutor, lo que vincula a los ejecutores con esos datos. No se pueden publicar hasta que todos los consumidores hayan terminado de leerlo. Ese acoplamiento es la principal razón por la que los clústeres no pueden reducirse rápidamente y por la que la pérdida de un ejecutor provoca costosos reintentos de etapa.

Una reducción de escala eficiente rompe este acoplamiento:

  • Las reorganizaciones de datos de gran tamaño van directamente a Azure Blob Storage a través del Administrador remoto de reorganización.
  • Las mezclas pequeñas permanecen en el disco local para mayor velocidad. Si ese ejecutor debe liberarse más adelante, Shuffle Migration mueve los bloques a pares o al almacenamiento de reserva en segundo plano.
  • La capa de decisión elige la ruta de acceso correcta por fase en tiempo de ejecución.
  • AQE Shuffle Write garantiza que el escritor genera particiones que el AQE de bajada consume sin volver a fusionar, evitando la E/S desperdiciada.
                ┌───────────────────────────┐
   Query  ───►  │   AQE + Decision Layer    │   per-stage choice
                └─────────────┬─────────────┘
                              │
                ┌─────────────▼─────────────┐
                │   AQE Shuffle Write       │   partition-aware writer
                └─────┬─────────────────┬───┘
                      │                 │
              local   ▼                 ▼   remote
        ┌────────────────────┐   ┌──────────────────┐
        │  Local disk +      │   │  RSM → Azure     │
        │  Shuffle Migration │   │  Blob Storage    │
        └─────────┬──────────┘   └─────────┬────────┘
                  │ on decommission        │
                  ▼                        ▼
        fallback storage   Remote shuffle store

Enrutamiento inteligente (capa de decisión)

La capa de decisión evalúa cada intercambio aleatorio y decide:

  • Grandes redistribuciones → Azure Blob Storage. Máxima reducción y máximo beneficio en tolerancia a errores.
  • Pequeñas reorganizaciones → disco local. No hay sobrecarga de E/S en la nube para transferencias diminutas. Si el ejecutor se retira posteriormente, Shuffle Migration toma el relevo.

El enrutamiento es automático y no requiere ninguna entrada de usuario. La granularidad recomendada es por fase.

Ventajas principales

Menores costes: paga solo por la capacidad de computación que utilices.

Con una reducción eficiente, los ejecutores se liberan en cuanto terminan su trabajo. Ya no permanecen inactivos almacenando datos de mezcla que las tareas posteriores podrían llegar a leer.

  • Reducción de escala más rápida. El escalado automático quita los nodos inmediatamente después de la finalización de la tarea.
  • Menos capacidad de cómputo ociosa. Ningún ejecutor "zombi" se mantuvo vivo sólo para servir a su orden aleatorio local.
  • Sin sobreaprovisionamiento de disco. Las reordenaciones de datos grandes van al almacenamiento de blobs en lugar de requerir discos locales grandes.
  • Costo de almacenamiento limitado. El almacenamiento alternativo se elimina automáticamente cuando los bloques dejan de ser necesarios.

Trabajos más resistentes

Cuando los datos de shuffle solo se almacenan en el disco local, el fallo de un ejecutor implica que esos datos se pierden y que Spark debe volver a calcularlos. Con una reducción vertical eficaz, los datos ya están en Blob Storage o se migran allí antes de que el ejecutor desaparezca.

Scenario Sin una reducción eficaz Con una reducción de escala eficaz
El ejecutor se bloquea Datos de shuffle perdidos; etapas reejecutadas Los datos son seguros en el almacenamiento; sin recomputación
Adelantamiento del nodo Datos desaparecidos, reintentos costosos Los datos se conservan; el trabajo continúa normalmente
Retirada ordenada La reproducción aleatoria se desactiva al apagar Bloques migrados al almacenamiento par o de reserva
Interrupciones de red durante la recuperación En cascada FetchFailedException Las lecturas proceden del almacenamiento y no se ven afectadas

Esto elimina la causa más común de FetchFailedException en producción.

Escalabilidad más rápida y realmente elástica

Sin una reducción de escala eficiente, el escalador automático no puede recuperar un nodo mientras algún ejecutor en él siga teniendo datos de shuffle o datos en caché. Una reducción de escala eficiente desacopla ambas cosas:

  • Los datos de intercambio están en el almacenamiento de blobs (o se migran allí al apagarse).
  • La caché ya no mantiene fijados los ejecutores. Las memorias caché reproducibles, como la caché de instantáneas Delta, están excluidas de la protección frente a la reducción de escala.

El escalador automático puede quitar libremente los nodos inactivos y cambiar el tamaño del clúster en respuesta a los cambios de carga de trabajo.

Mejor rendimiento con datos sesgados y shuffles grandes

AQE Shuffle Write permite a la Ejecución Adaptativa de Consultas configurar la propia escritura del shuffle, eligiendo una partición que el AQE posterior consume sin volver a fusionarla y generando menos bloques, con un tamaño más adecuado, para el almacenamiento remoto. En combinación con la capa de decisión, se obtiene un tiempo de reloj más rápido en consultas grandes o sesgadas y una latencia sin cambios para las pequeñas.

Empieza ahora

Aplique esta configuración para habilitar toda la solución eficiente de reducción de escala:

# Remote Shuffle Manager
spark.conf.set("spark.remote.shuffle.enabled", "true")

# Decision Layer — per-stage routing of local vs. remote shuffle
spark.conf.set("spark.sql.rsm.decisionlayer.enabled.level", "stage")

# AQE participates in shuffle write
spark.conf.set("spark.sql.adaptive.shuffleWrite.enabled", "true")

# Shuffle Migration on executor decommission
spark.conf.set("spark.storage.decommission.shuffleBlocks.enabled", "true")
spark.conf.set("spark.storage.decommission.shuffleBlocks.cleanup", "true")
spark.conf.set("spark.storage.decommission.shuffleBlocks.migrateToFallbackStorage", "true")
spark.conf.set("spark.storage.decommission.fallbackStorage.cleanUp", "true")

No se requiere ningún cambio de código. También puede establecerlos en las propiedades de Spark del entorno.

Referencia de configuración

Gestor remoto de shuffle (RSM)

Configuración Recomendado Qué controla
spark.remote.shuffle.enabled true Activa la reducción de escala eficiente. Los datos de shuffle se almacenan en Azure Blob Storage en lugar de en los discos locales del ejecutor.

Capa de decisión

Configuración Recomendado Qué controla
spark.sql.rsm.decisionlayer.enabled.level stage Nivel de granularidad en el que la capa de decisión dirige la mezcla. stage evalúa cada fase de Spark de forma independiente.

Escritura aleatoria de AQE

Configuración Recomendado Qué controla
spark.sql.adaptive.shuffleWrite.enabled true Permite que AQE participe en la fase de escritura aleatoria. Genera un particionado que AQE utiliza posteriormente sin tener que volver a compactarlo.

Note

El propio AQE (spark.sql.adaptive.enabled) debe estar activado. Está activado de forma predeterminada en Fabric Spark.

Migración aleatoria en la retirada

Configuración Recomendado Qué controla
spark.storage.decommission.shuffleBlocks.enabled true Migra los bloques de mezcla fuera de un ejecutor que se está retirando, en lugar de eliminarlos.
spark.storage.decommission.shuffleBlocks.cleanup true Limpia los bloques de shuffle del ejecutor de origen tras una migración correcta.
spark.storage.decommission.shuffleBlocks.migrateToFallbackStorage true Si ningún ejecutor del mismo nivel puede aceptar los bloques, los migra al almacenamiento de reserva (Azure Blob Storage).
spark.storage.decommission.fallbackStorage.cleanUp true Elimina los bloques de mezcla del almacenamiento de respaldo una vez que ya no son necesarios, acotando el coste de almacenamiento.

Asignación dinámica compatible con caché

Configuración Recomendado Qué controla
spark.dynamicAllocation.preventShutdownExecutorWithCache false La asignación dinámica permite liberar ejecutores incluso cuando mantienen bloques en caché.
spark.dynamicAllocation.excludeDeltaSnapshotCache true Omite la memoria caché de instantáneas delta al decidir si un ejecutor sigue manteniendo una memoria caché útil. La caché de instantáneas delta es reproducible y no debe bloquear la reducción vertical.

Optimización avanzada (RSM)

La mayoría de los usuarios no necesitan cambiar estos valores predeterminados.

Rendimiento de escritura

Configuración Predeterminado Qué controla
spark.remote.shuffle.partition.buffersize 16777216 (16 MB) Búfer por partición antes de escribir en el almacenamiento.
spark.remote.shuffle.blocksize 8388608 (8 MB) Tamaño de bloques individuales cargados en Blob Storage.
spark.remote.shuffle.write.maxthreads cores × 16 Número máximo de subprocesos utilizados para escribir datos de shuffle.
spark.remote.shuffle.write.maxtasks 16384 Número máximo de operaciones de escritura simultáneas.

Rendimiento de lectura

Configuración Predeterminado Qué controla
spark.remote.shuffle.read.parallel.enabled true Secuencias de descarga paralelas para lecturas aleatorias.
spark.remote.shuffle.read.parallelism 4 Secuencias de descarga paralelas por tarea.
spark.remote.shuffle.read.prefetchqueuesize 250 Captura previa de la profundidad de la cola durante las lecturas.
spark.remote.shuffle.read.maxthreads cores × 4 Número máximo de subprocesos utilizados para la lectura.

Confiabilidad

Configuración Predeterminado Qué controla
spark.remote.shuffle.retries 5 Reintentar cuando se produzcan errores transitorios de almacenamiento.
spark.remote.shuffle.retrydelayms 800 Intervalo de espera inicial entre reintentos.
spark.remote.shuffle.retrymaxdelayms 60000 Límite de retroceso.

Compression

Configuración Predeterminado Qué controla
spark.remote.shuffle.compression Utiliza spark.io.compression.codec Formato de compresión para datos de redistribución remotos (por ejemplo, lz4, zstd).

Resultados de rendimiento

Gráfico que muestra el ahorro en costes de computación con una reducción de escala eficiente habilitada frente a cuando está deshabilitada en una prueba de referencia TPC-DS, que demuestra una reducción de costes del 54 %.

Ahorro en costes de computación (benchmark TPC-DS)

Métrica Sin una reducción eficaz Con una reducción de escala eficaz
Cómputo total (minutos de VM) 14,952 6,880
Reducción de costos 54%

El tiempo total de ejecución del trabajo puede ser mayor (el escalado automático usa menos ejecutores simultáneos), pero la capacidad de cómputo facturada se reduce en más de la mitad.

Rendimiento de la capa de decisión (TPC-DS, RSM activado)

Dirigir los shuffles pequeños al disco local y solo los shuffles grandes al almacenamiento remoto ofrece una mejora del tiempo de ejecución de hasta un 57% frente a dirigir todos los shuffles al almacenamiento remoto, con el mismo beneficio de reducción de escala.

Limitaciones

  • Se requiere NEE. La reducción de escala eficiente depende del motor de ejecución nativo.
  • Solo Azure Blob Storage. Estándar BlockBlobStorage con HNS deshabilitado. Las cuentas de Azure Data Lake Gen2 o con HNS habilitado no son compatibles como almacén remoto de shuffle.
  • No se admite con Azure Private Link. Actualmente, los entornos que usan redes de private link no son compatibles.
  • La granularidad de la capa de decisión se encuentra actualmente por fase. El enrutamiento por tarea o por partición no se contempla.
  • Cambio en el comportamiento de la caché. Con preventShutdownExecutorWithCache=false, es posible que se reduzcan los ejecutores que contienen datos de cache()/persist(). Las cargas de trabajo que dependen en gran medida de la caché local del ejecutor para los datos activos deben validarse.