Optimisation des jointures de plages
Une jointure de plage se produit lorsque deux relations se joignent à l'aide d'une condition de point dans un intervalle ou de chevauchement d'intervalle. L'utilisation de l'optimisation de jointure de plage dans Databricks Runtime peut grandement améliorer les performances des queries.
Dans Databricks SQL, Databricks optimise automatiquement les jointures par plage sans aucune configuration manuelle. Vous pouvez également ajuster manuellement les jointures par plage à l'aide d'indicateurs de jointure ou de la configuration de session pour tous les types de compute.
Point dans la jointure de plage d’intervalle
Une jointure de plage d'intervalles ponctuels est une jointure dont la condition contient des prédicats spécifiant qu’une valeur d’une relation se situe entre deux valeurs de l’autre relation. Par exemple :
-- using BETWEEN expressions
SELECT *
FROM points JOIN ranges ON points.p BETWEEN ranges.start and ranges.end;
-- using inequality expressions
SELECT *
FROM points JOIN ranges ON points.p >= ranges.start AND points.p < ranges.end;
-- with fixed length interval
SELECT *
FROM points JOIN ranges ON points.p >= ranges.start AND points.p < ranges.start + 100;
-- join two sets of point values within a fixed distance from each other
SELECT *
FROM points1 p1 JOIN points2 p2 ON p1.p >= p2.p - 10 AND p1.p <= p2.p + 10;
-- a range condition together with other join conditions
SELECT *
FROM points, ranges
WHERE points.symbol = ranges.symbol
AND points.p >= ranges.start
AND points.p < ranges.end;
Jointure de plage de chevauchement d'intervalle
Un join de plages de chevauchement d'intervalle est un join dans lequel la condition contient des prédicats spécifiant un chevauchement d'intervalles entre deux valeurs de chaque relation. Par exemple :
-- overlap of [r1.start, r1.end] with [r2.start, r2.end]
SELECT *
FROM r1 JOIN r2 ON r1.start < r2.end AND r2.start < r1.end;
-- overlap of fixed length intervals
SELECT *
FROM r1 JOIN r2 ON r1.start < r2.start + 100 AND r2.start < r1.start + 100;
-- a range condition together with other join conditions
SELECT *
FROM r1 JOIN r2 ON r1.symbol = r2.symbol
AND r1.start <= r2.end
AND r1.end >= r2.start;
Optimisation des jointures de plages
L'optimisation de la jonction de plage est effectuée pour les jonctions qui :
- Avoir une condition qui peut être interprétée comme un point dans un intervalle ou une jointure de plages à chevauchement d'intervalles.
- Toutes les valeurs impliquées dans la condition de jointure d'étendue sont de type numérique (entier, à virgule flottante, décimal),
DATEouTIMESTAMP. - Toutes les valeurs impliquées dans la condition de jointure par plage sont du même type. Dans le cas du type décimal, les valeurs doivent également avoir la même échelle et la même précision.
- Il s'agit d'un(e)
INNER JOIN, ou dans le cas d'une jointure de plage d'intervalle par point, d'unLEFT OUTER JOINavec une valeur de point sur le côté gauche, ou d'unRIGHT OUTER JOINavec une valeur de point sur le côté droit. - Avoir une taille de compartiment, soit dérivée automatiquement, soit spécifiée manuellement.
Jointures avec égalité numérique et conditions de plage
Lorsqu'une condition de jointure inclut à la fois une condition d'égalité sur une colonne numérique et une condition de plage, l'optimiseur peut appliquer le binning à la colonne d'égalité numérique car elle satisfait aux exigences de type pour l'optimisation des jointures par plage. Cela peut entraîner l'affectation de la colonne d'égalité à des compartiments ou son exclusion de l'optimisation, ce qui réduit les performances.
Pour vous assurer que l'optimisation de la jointure de plage s'applique uniquement à la condition de plage prévue, convertissez les colonnes d'égalité numérique en STRING. Cela les exclut d'être considérées comme des colonnes de condition de plage.
SELECT /*+ RANGE_JOIN(reference, 3306084) */
reference.*, position.*
FROM position
INNER JOIN reference
ON CAST(position.parent_index AS STRING) = CAST(reference.parent_index AS STRING)
AND position.child_index BETWEEN reference.min_child_index AND reference.max_child_index;
Le même modèle s’applique aux autres colonnes numériques utilisées comme clés d’égalité, telles que DATE, les identifiants entiers ou les colonnes de partition clusterisées.
Taille du compartiment
La *taille du compartiment* est un paramètre de réglage numérique qui divise le domaine des valeurs de la condition de plage en plusieurs compartiments de taille égale. Par exemple, avec une taille de bin de 10, l'optimisation divise le domaine en bins qui sont des intervalles de longueur 10.
Si vous avez un point dans la condition de plage de p BETWEEN start AND end, et start est 8 et end est 22, cet intervalle de valeurs chevauche trois bins de longueur 10 – le 1er bin de 0 à 10, le 2e bin de 10 à 20 et le 3e bin de 20 à 30. Seuls les points qui tombent dans les trois mêmes bins doivent être considérés comme des correspondances de jointure possibles pour cet intervalle. Par exemple, si p est 32, il peut être exclu de se situer entre start de 8 et end de 22, car il se trouve dans l'intervalle de 30 à 40.
- Pour
DATEvaleurs, la valeur de la taille des intervalles est interprétée comme des jours. Par exemple, une valeur de taille de compartiment de 7 représente une semaine. - Pour
TIMESTAMPvaleurs, la valeur de la taille de compartiment est interprétée comme des secondes. Si une valeur inférieure à la seconde est requise, des valeurs fractionnaires peuvent être utilisées. Par exemple, une valeur de taille de compartiment de 60 représente une minute, et une valeur de taille de compartiment de 0,1 représente 100 millisecondes.
Vous pouvez spécifier la taille du compartiment en utilisant une indication de jointure d'intervalle dans la query ou en définissant un parameter de configuration de session. Dans Databricks SQL, la taille du compartiment est automatiquement dérivée lorsque l'optimisation automatique des jointures de plage est activée.
Optimisation automatique des jointures de plage
Dans Databricks SQL, Databricks détecte automatiquement les jointures de plages admissibles et déduit la taille de compartiment optimale en échantillonnant la table des intervalles. Cela supprime la nécessité de spécifier manuellement une taille de compartiment via des indicateurs ou la configuration de session.
L'optimisation automatique des jointures de plage est activée by default dans Databricks SQL. Pour le désactiver, définissez la configuration suivante :
SET spark.databricks.optimizer.autoRangeJoin.enabled = false;
Si vous spécifiez une taille de bin via un indice de jointure de plage ou une configuration de session, cette valeur remplace la taille de bin automatiquement dérivée.
Activer la jointure de plage en utilisant un conseil de jointure de plage
Pour activer l'optimisation des jointures de plages dans une requête SQL, utilisez une *indication de jointure de plages* pour spécifier la taille du compartiment. L'indicateur doit contenir le nom de la relation de l'une des relations jointes et le paramètre numérique de taille de compartiment. Le nom de la relation peut être une table, une vue ou une sous-requête.
SELECT /*+ RANGE_JOIN(points, 10) */ *
FROM points JOIN ranges ON points.p >= ranges.start AND points.p < ranges.end;
SELECT /*+ RANGE_JOIN(r1, 0.1) */ *
FROM (SELECT * FROM ranges WHERE ranges.amount < 100) r1, ranges r2
WHERE r1.start < r2.start + 100 AND r2.start < r1.start + 100;
SELECT /*+ RANGE_JOIN(c, 500) */ *
FROM a
JOIN b ON (a.b_key = b.id)
JOIN c ON (a.ts BETWEEN c.start_time AND c.end_time)
Dans le troisième exemple, vous devez placer l’indication sur c.
Ceci est dû au fait que les jointures sont associatives à gauche, de sorte que la query est interprétée comme (a JOIN b) JOIN c,
et l'indice sur a s'applique à la jointure de a avec b et non à la jointure avec c.
#create minute table
minutes = spark.createDataFrame(
[(0, 60), (60, 120)],
"minute_start: int, minute_end: int"
)
#create events table
events = spark.createDataFrame(
[(12, 33), (0, 120), (33, 72), (65, 178)],
"event_start: int, event_end: int"
)
#Range_Join with "hint" on the from table
(events.hint("range_join", 60)
.join(minutes,
on=[events.event_start < minutes.minute_end,
minutes.minute_start < events.event_end])
.orderBy(events.event_start,
events.event_end,
minutes.minute_start)
.show()
)
#Range_Join with "hint" on the join table
(events.join(minutes.hint("range_join", 60),
on=[events.event_start < minutes.minute_end,
minutes.minute_start < events.event_end])
.orderBy(events.event_start,
events.event_end,
minutes.minute_start)
.show()
)
Vous pouvez également placer une indication de jointure de plage sur l'un des DataFrames joints. Dans ce cas, l'indice contient uniquement le paramètre de taille de bac numérique.
val df1 = spark.table("ranges").as("left")
val df2 = spark.table("ranges").as("right")
val joined = df1.hint("range_join", 10)
.join(df2, $"left.type" === $"right.type" &&
$"left.end" > $"right.start" &&
$"left.start" < $"right.end")
val joined2 = df1
.join(df2.hint("range_join", 0.5), $"left.type" === $"right.type" &&
$"left.end" > $"right.start" &&
$"left.start" < $"right.end")
Activer la jointure de plages à l’aide de la configuration de session
Si vous ne souhaitez pas modifier la query, spécifiez la taille du compartiment en tant que paramètre de configuration.
SET spark.databricks.optimizer.rangeJoin.binSize=5
Ce paramètre de configuration s'applique à toute jonction avec une condition de plage. Cependant, une taille de compartiment différente définie via une indication de jointure de plage remplace toujours celle définie via le paramètre.
Choisissez la taille du compartiment
L'efficacité de l'optimisation des jointures de plages dépend du choix de la taille de bin appropriée.
Une petite taille de compartiment se traduit par un plus grand nombre de compartiments, ce qui facilite le filtrage des correspondances potentielles.
Cependant, cela devient inefficace si la taille du compartiment est nettement inférieure aux intervalles de valeur rencontrés, et que les intervalles de valeur chevauchent plusieurs intervalles de *compartiment*. Par exemple, avec une condition p BETWEEN start AND end, où start est 1 000 000 et end est 1 999 999, et une taille de compartiment de 10, l'intervalle de valeur chevauche 100 000 compartiments.
Si la durée de l'intervalle est assez uniforme et connue, nous vous recommandons de définir la taille du segment sur la durée typique attendue de l'intervalle de valeurs. Cependant, si la longueur de l'intervalle est variable et asymétrique, un équilibre doit être trouvé pour définir une taille de compartiment qui filtre efficacement les intervalles courts, tout en empêchant les intervalles longs de chevaucher trop de compartiments. En supposant une table ranges, avec des intervalles qui se situent entre les colonnes start et end, vous pouvez déterminer différents percentiles de la valeur de longueur d'intervalle asymétrique avec la query suivante :
SELECT
map_from_arrays(
ARRAY(0.5, 0.9, 0.99, 0.999, 0.9999),
APPROX_PERCENTILE(
end::DOUBLE - start::DOUBLE,
ARRAY(0.5, 0.9, 0.99, 0.999, 0.9999)
)
) AS bin_sizes
FROM
ranges;
Le cast de chaque colonne en DOUBLE avant la soustraction garantit le fonctionnement de la query, que les colonnes soient numériques, DATE ou TIMESTAMP valeurs.
Un paramètre recommandé de la taille du compartiment serait le maximum de la valeur au 90e centile, ou la valeur au 99e centile divisée par 10, ou la valeur au 99,9e centile divisée par 100, etc. La justification est :
- Si la valeur au 90e percentile correspond à la taille du compartiment, seuls 10 % des longueurs d'intervalle de valeur sont plus longues que l'intervalle du compartiment, elles couvrent donc plus de 2 intervalles de compartiments adjacents.
- Si la valeur au 99e centile correspond à la taille du bin, seulement 1 % des longueurs d'intervalle de valeur couvrent plus de 11 intervalles de bin adjacents.
- Si la valeur au 99,9e centile est la taille du bin, seulement 0,1 % des longueurs d'intervalle de valeur couvrent plus de 101 intervalles de bin adjacents.
- La même opération peut être répétée pour les valeurs aux 99,99e, 99,999e percentiles, etc., si nécessaire.
La méthode décrite limite la quantité d'intervalles de valeurs longues asymétriques qui se chevauchent avec plusieurs intervalles de bacs. La valeur de la taille de bin obtenue de cette manière n’est qu’un point de départ pour l’affinement ; les résultats réels pourraient dépendre de la charge de travail spécifique.