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 :

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

-
Pendant l'exécution

-
Après l'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 :

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

-
Coalescez dynamiquement les partitions : nœud
CustomShuffleReaderavec propriétéCoalesced

-
Gérer dynamiquement la jointure asymétrique : nœud
SortMergeJoinavec le champisSkewà true.

-
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.


Configuration
Dans cette section :
- Activer et désactiver l'exécution adaptative de query
- Activer le mélange à optimisation automatique
- Modifier dynamiquement la jointure de tri-Merge en jointure par hachage de diffusion
- Fusionner dynamiquement les partitions
- Gérer dynamiquement la jonction asymétrique
- Détecter et propager dynamiquement les relations vides.
Activer et désactiver l'exécution adaptative de query
Propriété |
|---|
spark.databricks.optimizer.adaptive.enabled Type : |
Activer le brassage auto-optimisé
Propriété |
|---|
spark.sql.shuffle.partitions Type : |
Modifier dynamiquement la jointure par sort Merge join en jointure de hachage par diffusion
Propriété |
|---|
spark.databricks.adaptive.autoBroadcastJoinThreshold Type : |
Coalescer dynamiquement les partitions
Propriété |
|---|
spark.sql.adaptive.coalescePartitions.enabled Type : |
spark.sql.adaptive.advisoryPartitionSizeInBytes Type : |
spark.sql.adaptive.coalescePartitions.minPartitionSize Type : |
spark.sql.adaptive.coalescePartitions.minPartitionNum Type : |
Gérer dynamiquement la jointure asymétrique
Propriété |
|---|
spark.sql.adaptive.skewJoin.enabled Type : |
spark.sql.adaptive.skewJoin.skewedPartitionFactor Type : |
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes Type : |
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 : |
Questions fréquemment posées (FAQ)
Dans cette section :
- Pourquoi l’AQE n’a-t-il pas diffusé une petite table de jointure ?
- Dois-je toujours utiliser un indice de stratégie de jointure de diffusion avec AQE activé ?
- Quelle est la différence entre l'indicateur de jonction asymétrique et l'optimisation de la jonction asymétrique d'AQE ? Lequel devrais-je utiliser ?
- Pourquoi AQE n'a-t-il pas ajusté automatiquement mon ordre de jointure ?
- Pourquoi l'AQE n'a-t-elle pas détecté mon déséquilibre de données ?
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 JOINne 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.