Les filtres d'exécution réduisent la durée des requêtes de jointure en élaguant la table de gauche avant l'exécution d'une jointure par hachage. Au lieu d'analyser chaque ligne, ApsaraDB for SelectDB génère un filtre au moment de l'exécution à partir de la table de droite et le propage vers la couche d'analyse. Ainsi, seules les lignes correspondantes participent à la jointure.
Les filtres d'exécution sont activés par défaut. ApsaraDB for SelectDB génère automatiquement des prédicats IN et des filtres de Bloom en fonction des statistiques de la requête et des tables. Utilisez des variables de session pour ajuster ce comportement pour des requêtes spécifiques.
Fonctionnement
Lors d'une jointure par hachage, la table de droite est chargée en premier pour construire une table de hachage. Un filtre d'exécution capture les valeurs clés de cette table de hachage et envoie le filtre au nœud OlapScanNode qui analyse la table de gauche. Le nœud OlapScanNode rejette alors les lignes non correspondantes avant qu'elles n'atteignent le nœud HashJoinNode.
L'exemple suivant utilise deux tables : T1 (une table de faits contenant 1 000 000 de lignes) et T2 (une table de dimensions contenant 200 lignes).
Sans filtre d'exécution, les 1 000 000 de lignes de T1 remontent jusqu'à la jointure :
HashJoinNode
| |
| 1,000,000 | 200
| |
OlapScanNode OlapScanNode
^ ^
| 1,000,000 | 200
T1 (fact) T2 (dimension)
Avec un filtre d'exécution appliqué au niveau de l'analyse, seules 6 000 lignes atteignent la jointure :
HashJoinNode
| |
| 6,000 | 200
| |
OlapScanNode OlapScanNode
^ ^
| 1,000,000 | 200
T1 (fact) T2 (dimension)
Lorsque le filtre est propagé plus bas vers le moteur de stockage, les index élaguent les données avant même leur lecture :
HashJoinNode
| |
| 6,000 | 200
| |
OlapScanNode OlapScanNode
^ ^
| 6,000 | 200
T1 (fact) T2 (dimension)
Contrairement à la propagation de prédicats ou à l'élagage de partitions, la condition de filtrage n'est pas connue lors de la planification de la requête. Elle est calculée à partir des données réelles de la table de droite pendant l'exécution, puis diffusée au nœud OlapScanNode qui lit la table de gauche.
Concepts clés
Table de gauche : table située du côté gauche d'une jointure, utilisée pour l'opération de sondage. La réorganisation des jointures peut modifier la position des tables.
Table de droite : table située du côté droit d'une jointure, utilisée pour construire la table de hachage. La réorganisation des jointures peut également modifier cette position.
Fragment : unité d'exécution de requête. Le nœud frontend (FE) divise une instruction SQL en fragments et les distribue aux nœuds backend (BE) du cluster distribué.
Cas d'utilisation des filtres d'exécution
Les filtres d'exécution sont particulièrement efficaces lorsque :
La table de gauche est nettement plus grande que la table de droite. La génération d'un filtre consomme de la mémoire et des ressources de calcul ; cette opération n'est rentable qu'à grande échelle.
Le résultat de la jointure est beaucoup plus petit que la table de gauche, ce qui signifie que le filtre peut éliminer la majorité des lignes de la table de gauche.
Si la table de droite est volumineuse ou si le résultat de la jointure est presque aussi grand que la table de gauche, les filtres d'exécution peuvent ajouter une surcharge sans apporter de bénéfice.
Types de filtres
ApsaraDB for SelectDB prend en charge cinq types de filtres d'exécution :
| Type | Fonctionnement | Idéal pour | Limitations |
|---|---|---|---|
| Prédicat IN | Construit un HashSet contenant toutes les valeurs clés de la table de droite et filtre la table de gauche avec IN |
Petites tables de droite dans les jointures par diffusion | Fonctionne uniquement pour les jointures par diffusion ; devient invalide lorsque le nombre de lignes de la table de droite dépasse runtime_filter_max_in_num (valeur par défaut : 1 024) |
| Filtre de Bloom | Construit une structure probabiliste à partir de la table de hachage ; présente un faible taux de faux positifs | Grandes tables de droite, la plupart des types de données | Surcharge de création et d'application plus élevée ; peut nuire aux performances si le taux de filtrage est faible ou si la table de gauche est petite |
| Filtre MinMax | Extrait la plage min/max de la table de droite et filtre les lignes situées en dehors de cette plage | Colonnes clés numériques avec des plages ne se chevauchant pas | Inefficace sur les colonnes non numériques (par exemple VARCHAR) ; aucun avantage lorsque les plages se chevauchent complètement |
| IN_OR_BLOOM_FILTER | Sélectionne automatiquement un prédicat IN ou un filtre de Bloom en fonction du nombre de lignes de la table de droite au moment de l'exécution | Usage général (par défaut) | Seuil contrôlé par runtime_filter_max_in_num ; utilise IN en dessous de 102 400 lignes, filtre de Bloom au-dessus |
| Filtre Bitmap | Filtre à l'aide d'une colonne bitmap renvoyée par une sous-requête IN | Requêtes avec sous-requêtes IN bitmap | Pris en charge uniquement par le moteur vectorisé |
Configurer les filtres d'exécution
Utilisez les variables de session suivantes pour ajuster le comportement des filtres d'exécution pour des requêtes spécifiques.
Paramètres
| Paramètre | Valeur par défaut | Description |
|---|---|---|
runtime_filter_mode |
GLOBAL |
Contrôle la portée de propagation des filtres entre les fragments d'exécution. Valeurs : OFF, LOCAL, GLOBAL. |
runtime_filter_type |
IN_OR_BLOOM_FILTER |
Type(s) de filtre(s) à générer. Valeurs : IN, BLOOM_FILTER, MIN_MAX, IN_OR_BLOOM_FILTER, BITMAP_FILTER. Combinez-les avec des virgules. |
runtime_filter_wait_time_ms |
1000 |
Temps d'attente maximal (ms) respecté par le nœud OlapScanNode avant de lancer une analyse. S'applique par filtre : trois filtres impliquent jusqu'à 3 000 ms au total. |
runtime_filters_max_num |
10 |
Nombre maximal de filtres de Bloom par requête. Les filtres de Bloom ayant la sélectivité la plus élevée sont conservés lorsque la limite est dépassée. |
runtime_bloom_filter_min_size |
1048576 (1 Mio) |
Taille minimale du filtre de Bloom en octets. |
runtime_bloom_filter_max_size |
16777216 (16 Mio) |
Taille maximale du filtre de Bloom en octets. |
runtime_bloom_filter_size |
2097152 (2 Mio) |
Taille par défaut du filtre de Bloom lorsque la cardinalité de la table de droite n'est pas disponible. |
runtime_filter_max_in_num |
1024 |
Seuil de nombre de lignes au-delà duquel aucun prédicat IN n'est généré. Contrôle également le basculement entre IN et filtre de Bloom en mode IN_OR_BLOOM_FILTER (seuil : 102 400). |
runtime_filter_mode
Contrôle la portée de propagation des filtres entre les fragments d'exécution de requête.
**
LOCAL** : le producteur du filtre (HashJoinNode) et le consommateur (OlapScanNode) doivent se trouver dans le même fragment. Adapté aux jointures par diffusion. Surcharge réduite.**
GLOBAL** : les filtres sont fusionnés et transférés entre les fragments via le réseau. Nécessaire pour les jointures par brassage (shuffle joins) où le producteur et le consommateur se trouvent dans des fragments différents.
Le mode GLOBAL couvre tous les scénarios gérés par LOCAL, ainsi que les jointures par brassage. Basculez vers LOCAL si la surcharge liée à la fusion et à la transmission des filtres dépasse les gains obtenus lors de l'analyse pour une jointure par brassage spécifique. Définissez la valeur sur OFF pour désactiver entièrement les filtres d'exécution.
Pour plus de détails sur la conception technique de la fusion de filtres entre fragments, consultez ISSUE 6116.
runtime_filter_type
Spécifiez un seul type ou combinez-en plusieurs :
-- By name (comma-separated, quoted)
SET runtime_filter_type = "BLOOM_FILTER,IN,MIN_MAX";
-- Equivalent numeric form (1=IN, 2=BLOOM_FILTER, 4=MIN_MAX)
SET runtime_filter_type = 7;
Comportement du prédicat IN :
Utilise une stratégie de fusion pour l'exécution distribuée.
Lorsque IN et d'autres types de filtres sont spécifiés simultanément et que la table de droite reste dans la limite définie par
runtime_filter_max_in_num, le système abandonne les autres filtres. En effet, un prédicat IN étant exact, les filtres supplémentaires n'apportent aucun avantage. Cette optimisation s'applique uniquement lorsque le producteur et le consommateur se trouvent dans le même fragment.
Comportement du filtre de Bloom :
Présente un taux de faux positifs non nul, ce qui signifie qu'il peut laisser passer certaines lignes non correspondantes. Les résultats restent corrects ; seule l'efficacité du filtrage est affectée.
Peut être propagé vers le moteur de stockage, mais uniquement pour les colonnes clés de la table de gauche. Sans propagation au niveau du stockage, les performances se dégradent souvent.
Évitez d'utiliser des filtres de Bloom sur des colonnes à faible sélectivité ou sur de petites tables de gauche.
Comportement du filtre MinMax :
Particulièrement efficace lorsque les plages de valeurs des tables de droite et de gauche ne se chevauchent pas (par exemple, le maximum de la table de droite est inférieur au minimum de la table de gauche, ou inversement).
Sur les colonnes numériques (INT, BIGINT, DOUBLE), le chevauchement des plages réduit son efficacité à zéro.
Sur les colonnes non numériques (VARCHAR), il dégrade généralement les performances.
Comportement de IN_OR_BLOOM_FILTER :
Utilise un prédicat IN lorsque le nombre de lignes de la table de droite est inférieur à 102 400 ; bascule vers un filtre de Bloom au-delà de ce seuil.
Ajustez le seuil à l'aide de
runtime_filter_max_in_num.
Comportement du filtre Bitmap :
Applicable uniquement lorsque la sous-requête IN renvoie une colonne bitmap.
Nécessite le moteur vectorisé.
runtime_filter_wait_time_ms
Le nœud OlapScanNode attend jusqu'à cette durée pour chaque filtre d'exécution assigné avant de commencer une analyse. Avec trois filtres, l'attente maximale correspond à trois fois cette valeur.
Les filtres qui arrivent dans la fenêtre d'attente sont propagés vers le moteur de stockage. Ceux qui arrivent après le début de l'analyse sont appliqués comme filtres d'expression sur les données déjà analysées. Bien qu'efficaces, ils sont moins performants qu'une propagation au niveau du stockage.
Conseils d'ajustement :
Clusters chargés exécutant des jointures longues : augmentez cette valeur afin que les filtres aient le temps d'être construits et transmis avant le début de l'analyse.
Clusters peu chargés exécutant de nombreuses petites requêtes : diminuez cette valeur pour éviter d'ajouter une latence inutile à des requêtes qui se terminent en quelques secondes.
Si un cluster est très sollicité et exécute de nombreuses requêtes gourmandes en ressources ou longues, prolongez la durée d'attente pour éviter que les requêtes complexes ne manquent des opportunités d'optimisation. À l'inverse, si la charge du cluster est légère et qu'il exécute de nombreuses petites requêtes ne durant que quelques secondes, réduisez la durée d'attente pour éviter d'ajouter une latence d'une seconde à chaque requête.
runtime_filters_max_num
Limite le nombre de filtres de Bloom par requête. Les prédicats IN et les filtres MinMax ne sont pas comptabilisés.
Lorsque la limite est dépassée, le système conserve les filtres de Bloom ayant la sélectivité la plus élevée, c'est-à-dire ceux susceptibles de filtrer le plus grand nombre de lignes :
Selectivity = HashJoinNode cardinality / HashJoinNode left child cardinality
La cardinalité estimée par le FE pouvant être imprécise, le classement basé sur la sélectivité peut ne pas refléter parfaitement l'efficacité réelle du filtrage.
N'ajustez ce paramètre que lors de l'optimisation de requêtes de jointure lentes entre de grandes tables.
Paramètres de taille du filtre de Bloom
Le FE calcule la longueur du filtre de Bloom lors de la planification de la requête. Tous les filtres de Bloom des nœuds HashJoinNode d'une même requête doivent avoir la même longueur pour pouvoir être fusionnés.
Si la cardinalité de la table de droite est disponible dans les statistiques, le FE estime la taille optimale et l'arrondit à la puissance de 2 supérieure.
Si la cardinalité n'est pas disponible, le FE utilise
runtime_bloom_filter_sizecomme valeur par défaut.runtime_bloom_filter_min_sizeetruntime_bloom_filter_max_sizebornent la taille finale, quelle que soit l'estimation.
Des filtres de Bloom plus grands traitent les colonnes à haute cardinalité avec plus de précision, mais consomment davantage de mémoire. Si la précision du filtrage est insuffisante pour une colonne à haute cardinalité (millions de valeurs distinctes), augmentez runtime_bloom_filter_size et testez les performances résultantes.
Ajustez la taille du filtre de Bloom par requête, et non globalement.
Vérifier l'application des filtres d'exécution
Exécutez EXPLAIN pour consulter le plan de requête et confirmer que les filtres sont générés et consommés sur les colonnes attendues.
Côté jointure (filtre généré) :
runtime filters: RF000[in] <- table.columnCôté analyse (filtre appliqué) :
runtime filters: RF000[in] -> table.column
Exemple :
CREATE TABLE test (t1 INT) DISTRIBUTED BY HASH (t1) BUCKETS 2;
INSERT INTO test VALUES (1), (2), (3), (4);
CREATE TABLE test2 (t2 INT) DISTRIBUTED BY HASH (t2) BUCKETS 2;
INSERT INTO test2 VALUES (3), (4), (5);
EXPLAIN SELECT t1 FROM test JOIN test2 WHERE test.t1 = test2.t2;
+-------------------------------------------------------------------+
| Explain String |
+-------------------------------------------------------------------+
| PLAN FRAGMENT 0 |
| OUTPUT EXPRS:`t1` |
| |
| 4:EXCHANGE |
| |
| PLAN FRAGMENT 1 |
| OUTPUT EXPRS: |
| PARTITION: HASH_PARTITIONED: `default_cluster:ssb`.`test`.`t1` |
| |
| 2:HASH JOIN |
| | join op: INNER JOIN (BUCKET_SHUFFLE) |
| | equal join conjunct: `test`.`t1` = `test2`.`t2` |
| | runtime filters: RF000[in] <- `test2`.`t2` |
| | |
| |----3:EXCHANGE |
| | |
| 0:OlapScanNode |
| TABLE: test |
| runtime filters: RF000[in] -> `test`.`t1` |
| |
| PLAN FRAGMENT 2 |
| OUTPUT EXPRS: |
| PARTITION: HASH_PARTITIONED: `default_cluster:ssb`.`test2`.`t2` |
| |
| 1:OlapScanNode |
| TABLE: test2 |
+-------------------------------------------------------------------+
Côté jointure : le nœud 2:HASH JOIN dans le PLAN FRAGMENT 1 génère un prédicat IN (RF000) à partir de test2.t2. Les valeurs ne sont connues qu'au moment de l'exécution.
Côté analyse : le nœud 0:OlapScanNode applique RF000 pour filtrer test.t1 avant que les lignes n'atteignent la jointure.
L'exécution de la requête renvoie [3, 4], soit uniquement les deux lignes où t1 = t2.
Vérifiez l'efficacité du filtre à l'aide d'un profil :
Activez le profilage, puis exécutez la requête :
SET enable_profile = true;
Dans le profil, recherchez la section RuntimeFilter sous OLAP_SCAN_NODE :
RuntimeFilter:in:
HasPushDownToEngine: true -- filter reached the storage engine
AWaitTimeCost: 0ns -- no wait; filter arrived before scan started
EffectTimeCost: 2.76ms -- time spent applying the filter
Et le résultat du filtrage :
RowsVectorPredFiltered: 9,320,008 -- rows discarded by the filter
VectorPredEvalTime: 364.39ms -- time spent on filter evaluation
Un nombre élevé pour RowsVectorPredFiltered confirme l'efficacité du filtre. La valeur HasPushDownToEngine: true confirme qu'il a bien atteint la couche de stockage.
Règles de planification
Les filtres d'exécution sont générés et appliqués selon les règles suivantes. Le non-respect de ces règles entraîne l'ignorance des filtres ou la production de résultats incorrects.
Règles de génération :
Les filtres sont générés uniquement pour les conditions d'égalité dans les clauses JOIN ON. L'égalité sûre pour NULL (
<=>) est exclue, car les valeurs NULL de la table de gauche pourraient être filtrées incorrectement.Le type de l'expression source ne peut pas être
HLLouBITMAP.Les expressions source et cible ne peuvent pas être des constantes.
Les expressions source et cible ne peuvent pas être identiques.
Les types des expressions source et cible doivent correspondre (les filtres de Bloom sont basés sur le hachage). Si les types diffèrent, le système tente de convertir l'expression cible vers le type source.
Les filtres provenant de
PlanNode.Conjunctsne sont pas propagés, car ils peuvent produire des résultats incorrects. Par exemple, lorsqu'une sous-requêteINest réécrite en jointure, le système stocke la condition JOIN générée automatiquement dansPlanNode.Conjuncts; y appliquer un filtre d'exécution risque de supprimer des lignes des résultats.
Règles de propagation :
Les filtres ne peuvent être propagés que vers un nœud OlapScanNode. Les autres types de nœuds d'analyse ne sont pas pris en charge.
Les filtres ne peuvent pas être propagés vers la table de gauche des jointures externes gauches, des jointures externes complètes ou des anti-jointures.
L'expression cible doit référencer une colonne existante dans la table de base d'origine.
L'expression cible ne peut pas contenir d'expressions de vérification de NULL telles que
COALESCE,IFNULLouCASE.
Règles de conduction de colonnes :
Si la clause JOIN ON contient
A.k = B.k AND B.k = C.k, le filtre pourC.kpeut être propagé uniquement versB.k, et non versA.k.Si la clause JOIN ON contient
A.a + B.b = C.cet queA.aest équivalent àB.a, alorsA.apeut être remplacé parB.aet le filtre propagé versB. SiA.aetB.ane sont pas équivalents, le filtre ne peut pas être propagé versB, car l'expression cible doit être liée à une seule table de gauche.