Gestion des requêtes volumineuses dans les workflows interactifs
Un défi avec les workflows de données interactifs est la gestion des grandes query. Cela inclut les queries qui génèrent trop de lignes de sortie, récupèrent de nombreuses partitions externes ou computent sur des ensembles de données extrêmement volumineux. Ces queries peuvent être extrêmement lentes, saturer les Ressources de compute et empêcher d'autres utilisateurs de partager la même ressource de compute.
Query Watchdog est un processus qui empêche les requêtes de monopoliser les ressources de compute en examinant les causes les plus courantes des requêtes volumineuses et en terminant les requêtes qui dépassent un threshold. Cet article décrit comment activer et configurer Query Watchdog.
Query Watchdog est activé pour tous les calculs multifonction créés à l'aide de l'interface utilisateur.
Exemple de query disruptive
Un analyste effectue des queries ad hoc dans un data warehouse just-in-time. L'analyste utilise un compute partagé à mise à l'échelle automatique qui permet à plusieurs utilisateurs d'utiliser un seul compute simultanément. Supposons qu’il existe deux tableaux contenant chacun un million de lignes.
import org.apache.spark.sql.functions._
spark.conf.set("spark.sql.shuffle.partitions", 10)
spark.range(1000000)
.withColumn("join_key", lit(" "))
.createOrReplaceTempView("table_x")
spark.range(1000000)
.withColumn("join_key", lit(" "))
.createOrReplaceTempView("table_y")
Ces tailles de table sont gérables dans Apache Spark. Cependant, chacun d'eux comprend une colonne join_key avec une chaîne vide dans chaque ligne. Cela peut se produire si les données ne sont pas parfaitement propres ou s'il existe une asymétrie de données significative où certaines clés sont plus répandues que d'autres. Ces clés de jointure vides sont bien plus répandues que toute autre valeur.
Dans le code suivant, l'analyste joint ces deux tables sur leurs clés, ce qui produit mille milliards de résultats , et tous sont produits sur un seul exécuteur (l'exécuteur qui obtient la clé " ") :
SELECT
id, count(id)
FROM
(SELECT
x.id
FROM
table_x x
JOIN
table_y y
on x.join_key = y.join_key)
GROUP BY id
Cette query semble être en cours d'exécution. Mais sans connaître les données, l'analyste constate qu'il ne reste « qu'une » seule tâche à exécuter au cours du job. La query ne se termine jamais, ce qui laisse l'analyste frustré et confus quant à la raison de son échec.
Dans ce cas, il n'y a qu'une seule clé de jointure problématique. D'autres fois, il peut y en avoir beaucoup plus.
Activer et configurer Query Watchdog
Pour activer et configurer Query Watchdog, les étapes suivantes sont requises.
- Activez Watchdog avec
spark.databricks.queryWatchdog.enabled. - Configurez le runtime de la tâche avec
spark.databricks.queryWatchdog.minTimeSecs. - Afficher la sortie avec
spark.databricks.queryWatchdog.minOutputRows. - Configurez le rapport de sortie avec
spark.databricks.queryWatchdog.outputRatioThreshold.
Pour éviter qu'une query ne génère trop de lignes de sortie pour le nombre de lignes d'entrée, vous pouvez activer Query Watchdog et configurer le nombre maximal de lignes de sortie comme un multiple du nombre de lignes d'entrée. Dans cet exemple, nous utilisons un rapport de 1000 (le default).
spark.conf.set("spark.databricks.queryWatchdog.enabled", true)
spark.conf.set("spark.databricks.queryWatchdog.outputRatioThreshold", 1000L)
Cette dernière configuration déclare qu'une tâche donnée ne doit jamais produire plus de 1000 fois le nombre de lignes d'entrée.
Le rapport de sortie est entièrement personnalisable. Nous vous recommandons de commencer plus bas et de voir quel threshold fonctionne bien pour vous et votre équipe. Une plage de 1 000 à 10 000 est un bon point de départ.
Non seulement Query Watchdog empêche les utilisateurs de monopoliser les ressources de compute pour les Jobs qui ne se termineront jamais, mais il permet également de gagner du temps en faisant échouer rapidement une query qui n'aurait jamais abouti. Par exemple, la query suivante échouera après plusieurs minutes car elle dépasse le ratio.
SELECT
z.id
join_key,
sum(z.id),
count(z.id)
FROM
(SELECT
x.id,
y.join_key
FROM
table_x x
JOIN
table_y y
on x.join_key = y.join_key) z
GROUP BY join_key, z.id
Voici ce que vous verriez :

Il est généralement suffisant d'activer Query Watchdog et de définir le ratio de threshold sortie/entrée, mais vous avez également la possibilité de définir deux propriétés supplémentaires : spark.databricks.queryWatchdog.minTimeSecs et spark.databricks.queryWatchdog.minOutputRows. Ces propriétés spécifient le temps minimum qu'une tâche donnée dans une query doit s'exécuter avant son annulation et le nombre minimum de lignes de sortie pour une tâche dans cette query.
Par exemple, vous pouvez définir minTimeSecs à une valeur plus élevée si vous voulez lui donner la possibilité de produire un grand nombre de lignes par tâche. De même, vous pouvez définir spark.databricks.queryWatchdog.minOutputRows à dix millions si vous souhaitez arrêter une query uniquement après qu'une tâche de cette query a produit dix millions de lignes. En dessous, la query réussit, même si le ratio sortie/entrée a été dépassé.
spark.conf.set("spark.databricks.queryWatchdog.minTimeSecs", 10L)
spark.conf.set("spark.databricks.queryWatchdog.minOutputRows", 100000L)
Si vous configurez Query Watchdog dans un Notebook, la configuration ne persiste pas lors des redémarrages de compute. Si vous souhaitez configurer Query Watchdog pour tous les utilisateurs d'un compute, nous vous recommandons d'utiliser une compute configuration.
Détecter les query sur des dataset extrêmement volumineux
Une autre query volumineuse typique peut analyser une grande quantité de données provenant de tables/datasets volumineux. L'opération d'analyse peut durer longtemps et saturer les ressources compute (même la lecture des métadonnées d'une grande table Hive peut prendre beaucoup de temps). Vous pouvez définir maxHivePartitions pour éviter de récupérer trop de partitions d'une grande table Hive. De même, vous pouvez également définir maxQueryTasks pour limiter les queries sur un dataset extrêmement volumineux.
spark.conf.set("spark.databricks.queryWatchdog.maxHivePartitions", 20000)
spark.conf.set("spark.databricks.queryWatchdog.maxQueryTasks", 20000)
Quand devriez-vous activer Query Watchdog ?
Query Watchdog devrait être activé pour les compute analytiques ad hoc où les analystes SQL et les data scientists partagent un même compute et un administrateur doit s'assurer que les queries « fonctionnent bien » les unes avec les autres.
Quand devriez-vous désactiver Query Watchdog ?
En général, nous ne conseillons pas d’annuler de manière anticipée les requêtes utilisées dans un scénario ETL, car il n’y a généralement pas d’humain dans la boucle pour corriger l’erreur. Nous vous recommandons de désactiver Query Watchdog pour tous les computes d’analytique, à l’exception des computes ad hoc.