Aller au contenu principal

Exécution adaptative de requêtes

L'exécution adaptative des requêtes (AQE) est une réoptimisation de requête qui se produit pendant l'exécution de la requête.

La motivation de la réoptimisation au runtime est que Databricks dispose des statistiques les plus récentes et précises à la fin d'un échange de brassage et de diffusion (appelé étape de requête dans AQE). En conséquence, Databricks peut opter pour une meilleure stratégie physique, choisir une taille et un nombre de partitions post-shuffle optimaux, ou effectuer des optimisations qui nécessitaient auparavant des indications, par exemple, la gestion des jointures asymétriques.

Ceci peut être très utile lorsque la collecte de statistiques n'est pas activée ou lorsque les statistiques sont obsolètes. Il est également utile dans les cas où les statistiques dérivées statiquement sont inexactes, comme au milieu d'une query complexe, ou après l'apparition d'une asymétrie des données.

Fonctionnalités

AQE est activé par default. Il possède 4 fonctionnalités majeures :

  • Modifie dynamiquement le tri Merge Join en un Broadcast Hash Join.
  • Agrège dynamiquement les partitions (combine de petites partitions en partitions de taille raisonnable) après l'échange de brassage. Les très petites tâches ont un throughput d'E/S plus faible et ont tendance à souffrir davantage du surcoût de planification et du surcoût de configuration des tâches. La combinaison de petites tâches permet d'économiser des Ressources et d'améliorer le throughput des clusters.
  • Gère dynamiquement l'asymétrie dans la jointure par tri-Merge et la jointure hachée par mélange en divisant (et en répliquant si nécessaire) les tâches asymétriques en tâches de taille à peu près égale.
  • Détecte et propage dynamiquement les relations vides.

Application

AQE s'applique à toutes les requêtes qui sont :

  • Non-streaming
  • Contenir au moins un échange (généralement lorsqu'il y a une jointure, une agrégation ou une fenêtre), une sous-requête, ou les deux.

Toutes les queries auxquelles l'AQE est appliquée ne sont pas nécessairement réoptimisées. La réoptimisation peut ou non aboutir à un plan de query différent de celui compilé statiquement. Pour déterminer si le plan d'une query a été modifié par AQE, consultez la section suivante : Plans de query.

Plans de query

Cette section explique comment vous pouvez examiner les plans de query de différentes manières.

Dans cette section :

Spark UI

AdaptiveSparkPlan nœud

Les query appliquées par l'AQE contiennent un ou plusieurs nœuds AdaptiveSparkPlan, généralement en tant que nœud racine de chaque query principale ou sous-query. Avant que la query ne s'exécute ou pendant son exécution, le drapeau isFinalPlan du nœud AdaptiveSparkPlan correspondant s'affiche comme false; une fois l'exécution de la query terminée, le drapeau isFinalPlan passe à true.

Plan en évolution

Le diagramme du plan de requête évolue à mesure que l'exécution progresse et reflète le plan le plus actuel en cours d'exécution. Les nœuds déjà exécutés (pour lesquels des métriques sont disponibles) ne changeront pas, mais ceux qui ne l'ont pas été peuvent changer au fil du temps en raison de réoptimisations.

Voici un exemple de diagramme de plan de query :

Diagramme du plan de query

DataFrame.explain()

AdaptiveSparkPlan nœud

Les requêtes AQE contiennent un ou plusieurs nœuds AdaptiveSparkPlan, généralement en tant que nœud racine de chaque requête principale ou sous-requête. Avant l'exécution de la query ou pendant son exécution, l'indicateur isFinalPlan du nœud AdaptiveSparkPlan correspondant s'affiche comme false; une fois l'exécution de la query terminée, l'indicateur isFinalPlan passe à true.

Plan actuel et initial

Sous chaque nœud AdaptiveSparkPlan, il y aura à la fois le plan initial (le plan avant d’appliquer les optimisations AQE) et le plan actuel ou final, selon que l'exécution est terminée ou non. Le plan actuel évoluera au fur et à mesure que l'exécution progresse.

Statistiques de Runtime

Chaque étape de brassage et de diffusion contient des statistiques de données.

Avant l'exécution ou pendant l'exécution de l'étape, les statistiques sont des estimations au moment de la compilation, et le flag isRuntime est false, par exemple : Statistics(sizeInBytes=1024.0 KiB, rowCount=4, isRuntime=false);

Une fois l'exécution de l'étape terminée, les statistiques sont celles collectées au moment de l'exécution, et l'indicateur isRuntime deviendra true, par exemple : Statistics(sizeInBytes=658.1 KiB, rowCount=2.81E+4, isRuntime=true)

Voici un exemple de DataFrame.explain :

  • Avant l'exécution

    Avant l'exécution

  • Pendant l'exécution

    Pendant l'exécution

  • Après l'exécution.

    Après exécution

SQL EXPLAIN

AdaptiveSparkPlan nœud

Les queries avec AQE contiennent un ou plusieurs nœuds AdaptiveSparkPlan, généralement comme nœud racine de chaque query principale ou sous-query.

Aucun plan actuel

Étant donné que SQL EXPLAIN n'exécute pas la query, le plan actuel est toujours le même que le plan initial et ne reflète pas ce qui serait finalement exécuté par AQE.

Voici un exemple d'explication SQL :

Explication SQL

Efficacité

Le plan de requête changera si une ou plusieurs optimisations AQE prennent effet. Les effets de ces optimisations AQE sont démontrés par la différence entre les plans actuel et final, et le plan initial et les nœuds de plan spécifiques dans les plans actuel et final.

  • Modifier dynamiquement la jointure par tri-Merge en jointure de hachage de diffusion : différents nœuds de jointure physiques entre le plan actuel/final et le plan initial

    Chaîne de stratégie de jonction

  • Coalescez dynamiquement les partitions : nœud CustomShuffleReader avec propriété Coalesced

    Lecteur de brassage personnalisé

    Chaîne du lecteur de brassage personnalisé

  • Gérer dynamiquement la jointure asymétrique : nœud SortMergeJoin avec le champ isSkew à true.

    Plan de jointure asymétrique

    Chaîne de jointure à asymétrie

  • Détectez et propagez dynamiquement les relations vides : une partie (ou l'intégralité) du plan est remplacée par le nœud LocalTableScan, le champ de relation étant vide.

    Analyse de table locale

    Chaîne d'analyse de table locale

Configuration

Dans cette section :

Activer et désactiver l'exécution adaptative de query

Propriété

spark.databricks.optimizer.adaptive.enabled Type : Boolean Activer ou désactiver l'exécution de requêtes adaptatives. Valeur par default : true

Propriété

spark.databricks.optimizer.adaptive.enabled Type : Boolean Activer ou désactiver l'exécution de requêtes adaptatives. Valeur par default : true

Activer le brassage auto-optimisé

Propriété

spark.sql.shuffle.partitions Type : Integer Le nombre de partitions par default à utiliser lors du brassage des données pour les jointures ou les agrégations. Définir la valeur auto active le mélange auto-optimisé, qui détermine automatiquement ce nombre en fonction du plan de query et de la taille des données d'entrée de la query. Note : pour Structured Streaming, cette configuration ne peut pas être modifiée entre les redémarrages de query à partir du même emplacement de point de contrôle. Valeur par default : 200

Propriété

spark.sql.shuffle.partitions Type : Integer Le nombre de partitions par default à utiliser lors du brassage des données pour les jointures ou les agrégations. Définir la valeur auto active le mélange auto-optimisé, qui détermine automatiquement ce nombre en fonction du plan de query et de la taille des données d'entrée de la query. Note : pour Structured Streaming, cette configuration ne peut pas être modifiée entre les redémarrages de query à partir du même emplacement de point de contrôle. Valeur par default : 200

Modifier dynamiquement la jointure par sort Merge join en jointure de hachage par diffusion

Propriété

spark.databricks.adaptive.autoBroadcastJoinThreshold Type : Byte String Le threshold pour trigger le basculement vers la jointure de diffusion à l'exécution. Valeur par default : 30MB

Propriété

spark.databricks.adaptive.autoBroadcastJoinThreshold Type : Byte String Le threshold pour trigger le basculement vers la jointure de diffusion à l'exécution. Valeur par default : 30MB

Coalescer dynamiquement les partitions

Propriété

spark.sql.adaptive.coalescePartitions.enabled Type : Boolean Activer ou désactiver la fusion des partitions. Valeur par default : true

spark.sql.adaptive.advisoryPartitionSizeInBytes Type : Byte String La taille cible après le regroupement. Les tailles des partitions regroupées seront proches de, mais pas supérieures à, cette taille cible. Valeur par default : 64MB

spark.sql.adaptive.coalescePartitions.minPartitionSize Type : Byte String La taille minimale des partitions après coalescence. Les tailles de partition fusionnées ne seront pas inférieures à cette taille. Valeur par default : 1MB

spark.sql.adaptive.coalescePartitions.minPartitionNum Type : Integer Le nombre minimum de partitions après fusion. Non recommandé, car le paramètre remplace explicitement spark.sql.adaptive.coalescePartitions.minPartitionSize. Default value: 2 fois le nombre de cœurs de cluster

Propriété

spark.sql.adaptive.coalescePartitions.enabled Type : Boolean Activer ou désactiver la fusion des partitions. Valeur par default : true

spark.sql.adaptive.advisoryPartitionSizeInBytes Type : Byte String La taille cible après le regroupement. Les tailles des partitions regroupées seront proches de, mais pas supérieures à, cette taille cible. Valeur par default : 64MB

spark.sql.adaptive.coalescePartitions.minPartitionSize Type : Byte String La taille minimale des partitions après coalescence. Les tailles de partition fusionnées ne seront pas inférieures à cette taille. Valeur par default : 1MB

spark.sql.adaptive.coalescePartitions.minPartitionNum Type : Integer Le nombre minimum de partitions après fusion. Non recommandé, car le paramètre remplace explicitement spark.sql.adaptive.coalescePartitions.minPartitionSize. Default value: 2 fois le nombre de cœurs de cluster

Gérer dynamiquement la jointure asymétrique

Propriété

spark.sql.adaptive.skewJoin.enabled Type : Boolean Activer ou désactiver la gestion des jointures asymétriques. Valeur par default : true

spark.sql.adaptive.skewJoin.skewedPartitionFactor Type : Integer Un facteur qui, lorsqu'il est multiplié par la taille médiane de la partition, contribue à déterminer si une partition est asymétrique. Valeur par default : 5

spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes Type : Byte String Un threshold qui contribue à déterminer si une partition est asymétrique. Valeur par default : 256MB

Propriété

spark.sql.adaptive.skewJoin.enabled Type : Boolean Activer ou désactiver la gestion des jointures asymétriques. Valeur par default : true

spark.sql.adaptive.skewJoin.skewedPartitionFactor Type : Integer Un facteur qui, lorsqu'il est multiplié par la taille médiane de la partition, contribue à déterminer si une partition est asymétrique. Valeur par default : 5

spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes Type : Byte String Un threshold qui contribue à déterminer si une partition est asymétrique. Valeur par default : 256MB

Une partition est considérée comme asymétrique lorsque (partition size > skewedPartitionFactor * median partition size) et (partition size > skewedPartitionThresholdInBytes) sont true.

Détecter et propager dynamiquement les relations vides

Propriété

spark.databricks.adaptive.emptyRelationPropagation.enabled Type : Boolean Activer ou désactiver la propagation dynamique des relations vides. Valeur par default : true

Propriété

spark.databricks.adaptive.emptyRelationPropagation.enabled Type : Boolean Activer ou désactiver la propagation dynamique des relations vides. Valeur par default : true

Questions fréquemment posées (FAQ)

Dans cette section :

Pourquoi l'AQE n'a-t-il pas diffusé une petite table de jointure ?

Si la taille de la relation censée être diffusée est inférieure à ce threshold mais n’est toujours pas diffusée :

  • Vérifiez le type de jointure. La diffusion n’est pas prise en charge pour certains types de jointure, par exemple, la relation de gauche d’une LEFT OUTER JOIN ne peut pas être diffusée.
  • Il se peut également que la relation contienne de nombreuses partitions vides, auquel cas la majorité des tâches peuvent se terminer rapidement avec une jointure par tri-Merge ou elle peut potentiellement être optimisée par la gestion des jointures asymétriques. AQE évite de transformer ces jointures de tri-Merge en jointures hachées de diffusion si le pourcentage de partitions non vides est inférieur à spark.sql.adaptive.nonEmptyPartitionRatioForBroadcastJoin.

Dois-je toujours utiliser une indication de stratégie de jointure de diffusion avec AQE activé ?

Oui. Une jointure de broadcast planifiée statiquement est généralement plus performante qu'une jointure planifiée dynamiquement par AQE, car AQE pourrait ne pas passer à la jointure de broadcast avant d'avoir effectué un brassage pour les deux côtés de la jointure (moment où les tailles réelles des relations sont obtenues). Ainsi, l'utilisation d'un indicateur de broadcast peut toujours être un bon choix si vous connaissez bien votre query. AQE respectera les query hints de la même manière que l'optimisation statique, mais peut toujours appliquer des optimisations dynamiques qui ne sont pas affectées par les hints.

Quelle est la différence entre l'indice de jointure asymétrique et l'optimisation de la jointure asymétrique AQE ? Lequel devrais-je utiliser ?

Il est recommandé de se fier à la gestion des jointures asymétriques AQE plutôt que d’utiliser l’indicateur de jointure asymétrique, car la jointure asymétrique AQE est entièrement automatique et, en général, plus performante que son homologue d’indicateur.

Pourquoi AQE n'a-t-il pas ajusté automatiquement l'ordre de mes jointures ?

Le réordonnancement dynamique des jointures ne fait pas partie de l'AQE.

Pourquoi AQE n'a-t-il pas détecté mon biais de données ?

Deux conditions de taille doivent être satisfaites pour qu'AQE détecte une partition comme une partition asymétrique :

  • La taille de la partition est supérieure à spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes (default 256 Mo)
  • La taille de la partition est supérieure à la taille médiane de toutes les partitions multipliée par le facteur de partition asymétrique spark.sql.adaptive.skewJoin.skewedPartitionFactor (default 5)

De plus, la prise en charge de la gestion des inclinaisons est limitée pour certains types de jointures : par exemple, dans LEFT OUTER JOIN, seule l’inclinaison du côté gauche peut être optimisée.

Ancien

Le terme « Adaptive Execution » existe depuis Spark 1.6, mais la nouvelle AQE de Spark 3.0 est fondamentalement différente. En termes de fonctionnalités, Spark 1.6 ne gère que la partie « regroupement dynamique des partitions ». En termes d’architecture technique, le nouveau AQE est un framework de planification et de replanification dynamique des requêtes basé sur des statistiques d’exécution, qui prend en charge une variété d’optimisations telles que celles que nous avons décrites dans cet article et peut être étendu pour permettre d’autres optimisations potentielles.