Optimiseur basé sur les coûts
Spark SQL peut utiliser un optimiseur basé sur les coûts (CBO) pour améliorer les plans de requêtes. Ceci est particulièrement utile pour les queries avec des jointures multiples. Pour que cela fonctionne, il est essentiel de collecter les statistiques de table et de colonne et de les maintenir à jour.
Collecter des statistiques
Pour profiter pleinement du CBO, il est important de collecter à la fois les statistiques de colonne et les statistiques de table . Vous pouvez utiliser la commande ANALYZE TABLE pour collecter manuellement des statistiques.
Pour maintenir les statistiques à jour, exécutez ANALYZE TABLE après avoir écrit dans la table.
Utilisation ANALYZE
L'optimisation prédictive exécute automatiquement ANALYZE, une commande de collecte de statistiques, sur les tables gérées par Unity Catalog. Databricks recommande d'activer l'optimisation prédictive pour toutes les tables gérées par Unity Catalog afin de simplifier la maintenance des données et de réduire les coûts de stockage. Voir ANALYZE TABLE … COMPUTE STATISTICS.
Vous pouvez supprimer les statistiques de l'optimiseur afin que ce dernier ne s'y fie plus, par exemple, pour atténuer un problème de planification causé par les statistiques. Voir ANALYZE TABLE... DROP STATISTICS.
Vérifier les plans de query
Il existe plusieurs façons de vérifier le plan de query.
commande EXPLAIN
Pour vérifier si le plan utilise des statistiques, utilisez la commande SQL EXPLAIN.
Si les statistiques sont manquantes, le plan de query pourrait ne pas être optimal.
== Optimized Logical Plan ==
Aggregate [s_store_sk], [s_store_sk, count(1) AS count(1)L], Statistics(sizeInBytes=20.0 B, rowCount=1, hints=none)
+- Project [s_store_sk], Statistics(sizeInBytes=18.5 MB, rowCount=1.62E+6, hints=none)
+- Join Inner, (d_date_sk = ss_sold_date_sk), Statistics(sizeInBytes=30.8 MB, rowCount=1.62E+6, hints=none)
:- Project [ss_sold_date_sk, s_store_sk], Statistics(sizeInBytes=39.1 GB, rowCount=2.63E+9, hints=none)
: +- Join Inner, (s_store_sk = ss_store_sk), Statistics(sizeInBytes=48.9 GB, rowCount=2.63E+9, hints=none)
: :- Project [ss_store_sk, ss_sold_date_sk], Statistics(sizeInBytes=39.1 GB, rowCount=2.63E+9, hints=none)
: : +- Filter (isnotnull(ss_store_sk) && isnotnull(ss_sold_date_sk)), Statistics(sizeInBytes=39.1 GB, rowCount=2.63E+9, hints=none)
: : +- Relation[ss_store_sk,ss_sold_date_sk] parquet, Statistics(sizeInBytes=134.6 GB, rowCount=2.88E+9, hints=none)
: +- Project [s_store_sk], Statistics(sizeInBytes=11.7 KB, rowCount=1.00E+3, hints=none)
: +- Filter isnotnull(s_store_sk), Statistics(sizeInBytes=11.7 KB, rowCount=1.00E+3, hints=none)
: +- Relation[s_store_sk] parquet, Statistics(sizeInBytes=88.0 KB, rowCount=1.00E+3, hints=none)
+- Project [d_date_sk], Statistics(sizeInBytes=12.0 B, rowCount=1, hints=none)
+- Filter ((((isnotnull(d_year) && isnotnull(d_date)) && (d_year = 2000)) && (d_date = 2000-12-31)) && isnotnull(d_date_sk)), Statistics(sizeInBytes=38.0 B, rowCount=1, hints=none)
+- Relation[d_date_sk,d_date,d_year] parquet, Statistics(sizeInBytes=1786.7 KB, rowCount=7.30E+4, hints=none)
La statistique rowCount est particulièrement importante pour les queries comportant plusieurs jointures. Si rowCount est manquant, cela signifie qu'il n'y a pas suffisamment d'informations pour le calculer (c'est-à-dire que certaines colonnes requises ne contiennent pas de statistiques).
Dans Databricks Runtime 16.0 et versions supérieures, la sortie de la commande EXPLAIN liste les tables référencées qui présentent des statistiques manquantes, partielles et complètes, comme dans l'exemple de sortie suivant :
== Optimizer Statistics (table names per statistics state) ==
missing = date_dim, store
partial =
full = store_sales
Corrective actions: consider running the following command on all tables with missing or partial statistics
ANALYZE TABLE <table-name> COMPUTE STATISTICS FOR ALL COLUMNS
Interface utilisateur Spark SQL
Utilisez la page de l'interface utilisateur Spark SQL pour voir le plan exécuté et la précision des statistiques.

Une ligne telle que rows output: 2,451,005 est: N/A signifie que cet opérateur produit environ 2 millions de lignes et qu'aucune statistique n'était disponible.

Une ligne telle que rows output: 2,451,005 est: 1616404 (1X) signifie que cet opérateur produit environ. 2 millions de lignes, alors que l'estimation était d'environ. 1,6 M et le facteur d'erreur d'estimation était de 1.

Une ligne telle que rows output: 2,451,005 est: 2626656323 signifie que cet opérateur produit environ 2M lignes alors que l'estimation était de 2B lignes, donc le facteur d'erreur d'estimation était de 1000.
Désactiver l'optimiseur basé sur les coûts
Le CBO est activé par default. Vous désactivez le CBO en modifiant l'indicateur spark.sql.cbo.enabled.
spark.conf.set("spark.sql.cbo.enabled", false)