Lors des requêtes de jointure multi-tables, les données ne correspondant pas à la condition de jointure augmentent la charge d'E/S et dégradent les performances. Dans ces scénarios, les runtime filters génèrent automatiquement des filtres légers pour éliminer les données non correspondantes dès la phase de scan. Cette approche réduit la surcharge d'E/S et améliore le temps de réponse des requêtes. Prise en charge pour les Hash Join depuis Hologres V2.0, cette fonctionnalité s'étend aux Cross Join depuis la version V4.2.
Contexte
Cas d'usage
Hologres prend en charge les runtime filters depuis la version V2.0. Cette fonctionnalité est couramment utilisée dans les scénarios de Hash Join impliquant au moins deux tables, en particulier lors de la jointure entre une grande table et une petite table. Aucune configuration manuelle n'est nécessaire. L'optimiseur et le moteur d'exécution optimisent automatiquement le filtrage des jointures au moment de la requête, ce qui diminue la charge d'E/S et accroît les performances.
À partir de la V4.2, les runtime filters couvrent également les scénarios de Cross Join. Ils optimisent spécifiquement le motif SQL courant où une sous-requête scalaire calcule un résultat agrégé ensuite utilisé comme condition de filtre sur une grande table. Pour plus d'informations, consultez Prise en charge des runtime filters pour les Cross Join (Nouveauté V4.2).
Fonctionnement
Runtime filters pour les hash join
Lors de la jointure de deux tables, les données de l'une sont chargées dans une table de hachage, tandis que les données de l'autre y sont comparées. Le processus de jointure comporte deux côtés :
Build side : côté qui construit la table de hachage, correspondant au nœud Hash dans le plan d'exécution.
Probe side : côté qui lit les données et les compare à la table de hachage du build side.
En règle générale, la petite table sert de build side et la grande table de probe side.
Les runtime filters fonctionnent en construisant un filtre léger à partir de la distribution des données du build side, puis en le propageant vers le probe side pour éliminer les données non pertinentes. Cela réduit le volume de données traitées par le probe side lors du Hash Join et minimise le trafic réseau, améliorant ainsi les performances de la jointure. Par conséquent, cette fonctionnalité est particulièrement efficace pour les jointures entre grandes et petites tables présentant une différence de taille significative, offrant des gains de performance supérieurs à ceux d'une jointure standard.
Runtime filters pour les cross join (Nouveauté V4.2)
Avant la V4.2, les runtime filters ne couvraient que les Hash Join. Cependant, il est fréquent d'utiliser des sous-requêtes scalaires pour calculer des valeurs agrégées, telles que min ou max, afin d'appliquer le résultat comme condition de filtre sur une grande table. Dans le plan d'exécution, ce motif SQL produit un Cross Join :
Build side : résultat de la sous-requête scalaire, composé d'exactement une ligne.
Probe side : scan de la grande table.
La V4.2 introduit un nouveau type de filtre, le ScalarFilter. L'optimiseur identifie automatiquement ce motif et propage la valeur scalaire évaluée du build side vers le ScanNode du probe side. Cela permet un filtrage au niveau des lignes et des RowGroups pendant la phase de scan, évitant ainsi un scan complet de la table suivi de comparaisons ligne par ligne.
Limitations et conditions de déclenchement
Limitations
Les runtime filters sont pris en charge uniquement dans Hologres V2.0 et versions ultérieures.
Pour les Hash Join, la V2.0 ne prend en charge les runtime filters que si la condition de jointure contient un seul champ. À partir de la V2.1, ils fonctionnent avec plusieurs champs.
Le TopN runtime filter est disponible uniquement à partir de la V4.0 et sert à améliorer les performances des calculs TopN sur table unique.
Le Cross Join runtime filter (ScalarFilter) nécessite la version V4.2 ou supérieure.
Conditions de déclenchement
Scénarios de Hash Join
Le moteur déclenche automatiquement un runtime filter lorsque toutes les conditions suivantes sont remplies :
Le probe side contient au moins 100 000 lignes.
Le ratio entre les données scannées du build side et celles du probe side est inférieur ou égal à 0,1. Plus ce ratio est faible, plus le déclenchement du filtre est probable.
Le ratio entre les données de sortie de la jointure et les données du probe side est inférieur ou égal à 0,1. Un ratio faible favorise également l'activation du filtre.
Scénarios de Cross Join (V4.2+)
L'optimiseur génère automatiquement un ScalarFilter lorsque les conditions ci-dessous sont satisfaites :
Le build side du Cross Join présente un nombre de lignes statistique égal à 1 et une distribution de type Replicated.
Les prédicats de filtre du probe side contiennent des conditions de comparaison avec des expressions issues du build side, telles que
>,<,>=,<=ouBETWEEN.
Types de runtime filters
Les runtime filters se classent selon les deux dimensions suivantes.
Par périmètre de shuffle (pour les hash join)
|
Type |
Versions prises en charge |
Scénarios |
|
Local |
V2.0+ |
Utilisé lorsque les données du probe side n'ont pas besoin d'être redistribuées (shuffled). Un runtime filter Local s'applique si les clés de jointure du build side et du probe side partagent la même distribution, si les données du build side sont diffusées (broadcast) vers le probe side, ou si elles sont redistribuées pour correspondre à la distribution du probe side. Ce type réduit uniquement le volume de données scannées et traitées par le Hash Join. |
|
Global |
V2.2+ |
Employé lorsque les données du probe side doivent être redistribuées. Le runtime filter est appliqué avant la redistribution des données, ce qui diminue le trafic réseau. |
Aucune spécification de type n'est requise. Le moteur effectue une sélection adaptative.
Par type de filtre
|
Type |
Versions prises en charge |
Description |
|
Bloom filter |
V2.0+ |
Filtre probabiliste susceptible de produire des faux positifs, ce qui signifie que certaines données pourraient ne pas être filtrées. Il reste néanmoins largement applicable et conserve une efficacité de filtrage élevée même lorsque le build side contient un grand volume de données. |
|
In filter |
V2.0+ |
Recommandé lorsque le build side présente un NDV (nombre de valeurs distinctes) faible. Il construit un HashSet à partir des données du build side et l'envoie au probe side pour filtrage. Ce filtre élimine précisément toutes les données requises et peut être combiné avec un index bitmap. |
|
MinMax filter |
V2.0+ |
Transmet les valeurs minimale et maximale des données du build side au probe side pour le filtrage. Il exploite les métadonnées pour ignorer des fichiers entiers ou des lots de données, réduisant ainsi les coûts d'E/S. |
|
ScalarFilter |
V4.2+ |
Conçu spécifiquement pour les scénarios de Cross Join. Lorsque le build side contient exactement une ligne, le moteur propage la valeur scalaire vers le ScanNode du probe side pour effectuer le filtrage pendant la phase de scan. |
Il est inutile de spécifier le type de filtre. Hologres sélectionne adaptativement le type approprié en fonction des conditions de jointure au moment de l'exécution.
Prise en charge des runtime filters pour les cross join (Nouveauté V4.2)
Motivation
Les utilisateurs calculent souvent des valeurs agrégées via des sous-requêtes scalaires pour utiliser les résultats comme conditions de filtre sur une grande table. Avant la V4.2, ce motif SQL présentait les problèmes suivants dans le plan d'exécution :
Le ScanNode du probe side ne pouvait pas exploiter la valeur du build side pour un filtrage anticipé et devait effectuer un scan complet de la table.
Après le scan, le Cross Join réalisait une comparaison ligne par ligne, entraînant un gaspillage important d'E/S.
Le runtime filter existant ne couvrait que les Hash Join et ne prenait pas en charge les Cross Join.
La V4.2 introduit le ScalarFilter pour résoudre ce problème. Cette fonctionnalité est activée par défaut et ne nécessite aucune intervention de l'utilisateur.
Comparaison SQL typique et plan d'exécution
SQL typique :
-- t1 is a large table; the aggregated min/max result from t2 is exactly 1 row.
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1);
Avant optimisation (sans runtime filter) :
Cross Join
-> Seq Scan on t1 -- Full table scan, no early filtering
-> Aggregate -- Build side min/max, 1-row result
-> Seq Scan on t2
Le probe side doit scanner toutes les données, puis le Cross Join compare les lignes une par une, ce qui provoque un gaspillage significatif d'E/S.
Après optimisation (V4.2 avec ScalarFilter) :
Cross Join
Runtime Filter Build Expr: (min(t2.a)), (max(t2.a))
-> Seq Scan on t1 -- Receives ScalarFilter, filters non-matching rows and RowGroups during scan
Runtime Filter Target Expr: (t1.a >= ${1}) AND (t1.a <= ${2})
-> Aggregate
-> Seq Scan on t2
Une fois le build side évalué, les valeurs réelles sont propagées vers le ScanNode du probe side, permettant d'effectuer le filtrage directement pendant la phase de scan.
Autres scénarios typiques
-- Scenario 1: Single table + scalar subquery
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1);
-- Scenario 2: Multi-table + scalar subquery
SELECT * FROM t1, t3
WHERE (SELECT min(a) FROM t2) <= t1.a
AND t3.a <= (SELECT max(a) FROM t2);
-- Scenario 3: CTE + Cross Join
WITH r AS (SELECT MIN(a) AS lo, MAX(a) AS hi FROM t2)
SELECT t1.* FROM t1, r
WHERE t1.a >= r.lo AND t1.a <= r.hi;
Fonctionnement
Identification par l'optimiseur : lorsque le build side d'un Cross Join affiche un nombre de lignes statistique égal à 1 et une distribution Replicated, l'optimiseur extrait les expressions de comparaison des prédicats (en développant
BETWEENen>=et<=) et génère des candidats ScalarFilter.Propagation au moment de l'exécution : durant la phase Open, le Cross Join consomme les données du build side, évalue l'expression
build_exprpour produire une valeur scalaire, construit un ScalarFilter et le publie vers le ScanNode du probe side.-
Filtrage anticipé dans le ScanNode : après réception du ScalarFilter par le NiagaraScan du probe side, celui-ci remplace les espaces réservés dans l'expression
target_exprpar les valeurs réelles et génère deux niveaux de filtrage :Filtrage au niveau de la ligne (évaluation conjonctive).
Filtrage au niveau du RowGroup (le moteur de stockage ignore les blocs de données non correspondants).
Limitations
Le build side doit contenir exactement une ligne. L'optimiseur ne génère un ScalarFilter que si le build side présente un nombre de lignes statistique égal à 1 et une distribution Replicated. Si le build side ne produit pas exactement une ligne lors de l'exécution, Hologres signale une erreur
RT_CHECK.Si la valeur du build side est NULL, Hologres publie un filtre
FILTER_ALL(qui filtre toutes les lignes), efface les données du build side pour court-circuiter le Cross Join et retourne un ensemble de résultats vide.Les colonnes dictionnaires ne sont pas prises en charge. Si la colonne du probe side référencée par
target_exprest une colonne encodée par dictionnaire (type DICTIONARY), le ScalarFilter n'est pas appliqué.Les types de données pris en charge incluent
INT8,UINT8,INT16,UINT16,INT32,UINT32,INT64,UINT64,DATE32,TIMESTAMP,FLOAT,DOUBLEetSTRING. Si un type de données n'est pas pris en charge, le ScalarFilter n'est pas propagé et aucune erreur n'est signalée.Les opérateurs de comparaison pris en charge sont
>,<,>=et<=. L'opérateurBETWEENest développé en deux comparaisons de plage.
Paramètres GUC
|
Paramètre |
Type |
Valeur par défaut |
Niveau |
Description |
|
|
bool |
|
|
Contrôle la génération d'un ScalarFilter pour les Cross Join. Cette fonctionnalité est activée par défaut et ne nécessite aucune action de l'utilisateur. |
Pour désactiver cette fonctionnalité, exécutez la commande suivante :
SET hg_experimental_generate_runtime_scalar_filter = off;
Nouveaux champs dans la sortie EXPLAIN
Lorsque le ScalarFilter est activé, le nœud Cross Join affiche les informations suivantes :
Runtime Filter Build Expr: expression du build side, par exemple(min(t2.a))et(max(t2.a)).Runtime Filter Target Expr: expression de filtre cible sur le probe side, telle que(t1.a >= ${1}) AND (t1.a <= ${2}). Dans cette expression,${N}est un espace réservé remplacé au moment de l'exécution par la valeur du build side correspondant aufilter_id.
Avantages en termes de performances
Dans le jeu de données TPC-DS de 10 To, le motif SQL suivant est fréquent :
DELETE FROM inventory
WHERE inv_date_sk >= (SELECT min(d_date_sk) FROM date_dim WHERE d_date BETWEEN 'INV_S_1' AND 'INV_E_1')
AND inv_date_sk <= (SELECT max(d_date_sk) FROM date_dim WHERE d_date BETWEEN 'INV_S_1' AND 'INV_E_1');
Grâce à l'optimisation ScalarFilter de la V4.2, le temps de requête passe de 2,672 secondes à 0,552 seconde, soit une amélioration des performances d'environ 4,8x.
Vérification des runtime filters
Les exemples suivants illustrent comment vérifier les effets des runtime filters dans différents scénarios.
Exemple 1 : Condition de jointure sur colonne unique (type local)
BEGIN;
CREATE TABLE test1 (x int, y int);
CALL set_table_property('test1', 'distribution_key', 'x');
CREATE TABLE test2 (x int, y int);
CALL set_table_property('test2', 'distribution_key', 'x');
END;
INSERT INTO test1 SELECT t, t FROM generate_series(1, 100000) t;
INSERT INTO test2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;
Plan d'exécution :
QUERY PLAN
Gather (cost=0.00..10.20 rows=1000 width=16)
[40:1 id=100002 dop=1 time=9/9/9ms rows=1000(1000/1000/1000) mem=16/16/16KB open=2/2/2ms get_next=7/7/7ms]
-> Hash Join (cost=0.00..10.16 rows=1000 width=16)
Hash Cond: (test1.x = test2.x)
Runtime Filter Cond: (test1.x = test2.x)
[id=8 dop=40 time=9/4/2ms rows=1000(39/25/17) mem=6/5/5KB open=4/1/0ms get_next=6/2/0ms]
-> Local Gather (cost=0.00..5.11 rows=1000000 width=8)
[id=3 dop=40 time=6/1/0ms rows=1000(39/25/17) mem=600/600/600B open=1/0/0ms get_next=6/1/0ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..5.10 rows=1000000 width=8)
Runtime Filter Target Expr: test1.x
[id=2 split_count=40 time=11/7/7ms rows=1000(39/25/17) mem=41/41/41KB open=11/7/7ms get_next=0/0/0ms scan_rows=1000000 25270/25000/24697))]
-> Hash (cost=5.00..5.00 rows=1000 width=8)
[id=7 dop=40 time=3/1/0ms rows=1000(39/25/17) mem=396/396/396KB open=3/1/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
-> Local Gather (cost=0.00..5.00 rows=1000 width=8)
[id=5 dop=40 time=1/0/0ms rows=1000(39/25/17) mem=0/0/0B open=1/0/0ms get_next=1/0/0ms local_dop=0/0/0]
-> Seq Scan on test2 (cost=0.00..5.00 rows=1000 width=8)
[id=4 split_count=40 time=1/0/0ms rows=1000(39/25/17) mem=528/528/528B open=1/0/0ms get_next=1/0/0ms scan_rows=1000(39/25/17)]
La table
test2contient 1 000 lignes et la tabletest1en contient 100 000. Le ratio de taille des données entre le build side et le probe side est de 0,01 (inférieur à 0,1), ce qui satisfait la condition de déclenchement par défaut des runtime filters.Le scan du probe side sur
test1afficheRuntime Filter Target Expr, indiquant qu'un runtime filter a été propagé.Du côté du probe side,
scan_rowsvaut 100 000 (nombre de lignes lues depuis le stockage), tandis querowsvaut 1 000 (nombre de lignes après filtrage). Cet écart démontre l'efficacité du filtrage.
Exemple 2 : Condition de jointure multi-colonnes (V2.1+, type local)
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x int, y int);
CREATE TABLE test2 (x int, y int);
END;
INSERT INTO test1 SELECT t, t FROM generate_series(1, 1000000) t;
INSERT INTO test2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x AND test1.y = test2.y;
Plan d'exécution :
QUERY PLAN
Gather (cost=0.00..10.46 rows=1000 width=16)
[40:1 id=100003 dop=1 time=6/6/6ms rows=1000(1000/1000/1000) mem=600/600/600B open=0/0/0ms get_next=6/6/6ms]
-> Hash Join (cost=0.00..10.43 rows=1000 width=16)
Hash Cond: ((test1.x = test2.x) AND (test1.y = test2.y))
Runtime Filter Cond: ((test1.x = test2.x) AND (test1.y = test2.y))
[id=8 dop=40 time=5/3/3ms rows=1000(1000/25/0) mem=40/2/1KB open=1/0/0ms get_next=4/3/3ms]
-> Local Gather (cost=0.00..5.11 rows=1000000 width=8)
[id=5 dop=40 time=4/3/3ms rows=1000(1000/25/0) mem=600/600/600B open=0/0/0ms get_next=4/3/3ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..5.10 rows=1000000 width=8)
Runtime Filter Target Expr: (test1.x AND test1.y)
[id=4 split_count=40 time=7/5/5ms rows=1000(1000/25/0) mem=49/11/9KB open=7/5/5ms get_next=1/0/0ms scan_rows=1000000(32768/25000/24576)]
-> Hash (cost=5.02..5.02 rows=40000 width=8)
[id=7 dop=40 time=1/0/0ms rows=40000(1000/1000/1000) mem=417/417/417KB open=1/0/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
-> Broadcast (cost=0.00..5.02 rows=40000 width=8)
[40:40 id=100002 dop=40 time=1/0/0ms rows=40000(1000/1000/1000) mem=0/0/0B open=1/0/0ms get_next=0/0/0ms * ]
-> Local Gather (cost=0.00..5.00 rows=1000 width=8)
[id=3 dop=40 time=1/0/0ms rows=1000(1000/25/0) mem=600/600/600B open=1/0/0ms get_next=0/0/0ms local_dop=1/1/1]
-> Seq Scan on test2 (cost=0.00..5.00 rows=1000 width=8)
[id=2 split_count=40 time=1/0/0ms rows=1000(1000/25/0) mem=528/528/528B open=0/0/0ms get_next=1/0/0ms scan_rows=1000(1000/1000/1000)]
La condition de jointure porte sur plusieurs colonnes, et le runtime filter est également généré pour plusieurs colonnes.
Les données du build side étant diffusées (broadcast), un runtime filter Local est utilisé.
Exemple 3 : Type Global (V2.2+, shuffle join)
SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x int, y int);
CREATE TABLE test2 (x int, y int);
END;
INSERT INTO test1 SELECT t, t FROM generate_series(1, 100000) t;
INSERT INTO test2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;
Plan d'exécution :
QUERY PLAN
-> Hash Join (cost=0.00..10.08 rows=1000 width=16)
Hash Cond: (test1.x = test2.x)
Runtime Filter Cond: (test1.x = test2.x)
[id=9 dop=40 time=10/8/8ms rows=1000(34/25/13) mem=6/6/5KB open=2/1/1ms get_next=8/7/7ms]
-> Redistribution (cost=0.00..5.07 rows=100000 width=8)
Hash Key: test1.x
[40:40 id=100002 dop=40 time=8/7/7ms rows=1289(46/32/20) mem=512/432/0B open=0/0/0ms get_next=8/7/7ms * ]
-> Local Gather (cost=0.00..5.01 rows=100000 width=8)
[id=3 dop=40 time=9/2/0ms rows=1289(1042/32/0) mem=600/600/600B open=0/0/0ms get_next=9/2/0ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..5.01 rows=100000 width=8)
Runtime Filter Target Expr: test1.x
[id=2 split_count=40 time=11/3/0ms rows=1289(1042/32/0) mem=50512/11849/528B open=11/3/0ms get_next=1/0/0ms scan_rows=100000(8192/7692/1696)]
-> Hash (cost=5.00..5.00 rows=1000 width=8)
[id=8 dop=40 time=2/1/1ms rows=1000(34/25/13) mem=396/396/396KB open=2/1/1ms get_next=0/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
-> Redistribution (cost=0.00..5.00 rows=1000 width=8)
Hash Key: test2.x
[40:40 id=100003 dop=40 time=2/1/1ms rows=1000(34/25/13) mem=0/0/0B open=0/0/0ms get_next=2/1/1ms * ]
-> Local Gather (cost=0.00..5.00 rows=1000 width=8)
[id=5 dop=40 time=1/0/0ms rows=1000(1000/25/0) mem=600/600/600B open=1/0/0ms get_next=1/0/0ms local_dop=1/1/1]
-> Seq Scan on test2 (cost=0.00..5.00 rows=1000 width=8)
[id=4 split_count=40 time=0/0/0ms rows=1000(1000/25/0) mem=528/528/528B open=0/0/0ms get_next=0/0/0ms scan_rows=1000(1000/1000/1000)]
Les données du probe side sont redistribuées vers l'opérateur Hash Join. Le moteur utilise automatiquement un runtime filter Global pour accélérer la requête.
Exemple 4 : In filter avec index bitmap (V2.2+)
SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x text, y text);
CALL set_table_property('test1', 'distribution_key', 'x');
CALL set_table_property('test1', 'bitmap_columns', 'x');
CALL set_table_property('test1', 'dictionary_encoding_columns', '');
CREATE TABLE test2 (x text, y text);
CALL set_table_property('test2', 'distribution_key', 'x');
END;
INSERT INTO test1 SELECT t::text, t::text FROM generate_series(1, 10000000) t;
INSERT INTO test2 SELECT t::text, t::text FROM generate_series(1, 50) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;
Plan d'exécution :
QUERY PLAN
Gather (cost=0.00..11.70 rows=50 width=14)
[40:1 id=100002 dop=1 time=16/16/16ms rows=50(50/50/50) mem=2/2/2KB open=0/0ms get_next=16/16/16ms]
-> Hash Join (cost=0.00..11.70 rows=50 width=14)
Hash Cond: (test1.x = test2.x)
Runtime Filter Cond: (test1.x = test2.x)
[id=7 dop=40 time=15/9/3ms rows=50(3/1/0) mem=5132/3774/264B open=1/0/0ms get_next=14/8/3ms]
-> Local Gather (cost=0.00..6.26 rows=10000000 width=12)
[id=3 dop=40 time=14/8/3ms rows=50(3/1/0) mem=600/600/600B open=1/0/0ms get_next=14/8/3ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..6.06 rows=10000000 width=12)
Runtime Filter Target Expr: test1.x
[id=2 split_count=40 time=16/10/5ms rows=50(3/1/0) mem=67544/48945/528B open=16/9/5ms get_next=1/0/0ms scan_rows=7247692(250875/249920/248982) bitmap_used=50]
-> Hash (cost=5.00..5.00 rows=50 width=2)
[id=6 dop=40 time=1/0/0ms rows=61(3/1/1) mem=534/530/521KB open=1/0/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=512/512/512KB]
-> Local Gather (cost=0.00..5.00 rows=50 width=2)
[id=5 dop=40 time=1/0/0ms rows=50(3/1/0) mem=600/600/600B open=0/0/0ms get_next=1/0/0ms local_dop=1/1/1]
-> Seq Scan on test2 (cost=0.00..5.00 rows=50 width=2)
[id=4 split_count=40 time=1/0/0ms rows=50(3/1/0) mem=528/528/528B open=1/0/0ms get_next=0/0/0ms scan_rows=50(3/1/1)]
L'opérateur de scan du probe side utilise un index bitmap. Le In filter assure un filtrage précis, ne laissant que 50 lignes. La valeur scan_rows dans l'opérateur de scan dépasse 7 millions, ce qui est inférieur aux 10 millions de lignes initiales. En effet, le In filter peut être propagé vers le moteur de stockage, réduisant ainsi la charge d'E/S. La combinaison d'un In filter avec un index bitmap offre des gains significatifs lorsque la clé de jointure est de type STRING.
Exemple 5 : MinMax filter pour la réduction des E/S (V2.2+)
SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x int, y int);
CALL set_table_property('test1', 'distribution_key', 'x');
CREATE TABLE test2 (x int, y int);
CALL set_table_property('test2', 'distribution_key', 'x');
END;
INSERT INTO test1 SELECT t::int, t::int FROM generate_series(1, 10000000) t;
INSERT INTO test2 SELECT t::int, t::int FROM generate_series(1, 100000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;
Plan d'exécution :
QUERY PLAN
Gather (cost=0.00..15.68 rows=100000 width=16)
[40:1 id=100002 dop=1 time=5/5/5ms rows=100000(100000/100000/100000) mem=600/600/600B open=0/0/0ms get_next=5/5/5ms]
-> Hash Join (cost=0.00..11.98 rows=100000 width=16)
Hash Cond: (test1.x = test2.x)
Runtime Filter Cond: (test1.x = test2.x)
[id=7 dop=40 time=5/4/4ms rows=100000(2639/2500/2406) mem=97/92/89KB open=1/0/0ms get_next=4/3/3ms]
-> Local Gather (cost=0.00..6.14 rows=10000000 width=8)
[id=3 dop=40 time=5/3/3ms rows=100000(2639/2500/2406) mem=600/600/600B open=1/0/0ms get_next=4/3/3ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..6.00 rows=10000000 width=8)
Runtime Filter Target Expr: test1.x
[id=2 split_count=40 time=6/6/5ms rows=100000(2639/2500/2406) mem=61/60/59KB open=6/5/5ms get_next=0/0/0ms scan_rows=327680(8192/8192/8192)]
-> Hash (cost=5.01..5.01 rows=100000 width=8)
[id=6 dop=40 time=1/0/0ms rows=100000(2639/2500/2406) mem=463/460/458KB open=1/0/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
-> Local Gather (cost=0.00..5.01 rows=100000 width=8)
[id=5 dop=40 time=1/0/0ms rows=100000(2639/2500/2406) mem=600/600/600B open=0/0/0ms get_next=1/0/0ms local_dop=1/1/1]
-> Seq Scan on test2 (cost=0.00..5.01 rows=100000 width=8)
[id=4 split_count=40 time=1/0/0ms rows=100000(2639/2500/2406) mem=528/528/528B open=1/0/0ms get_next=0/0/0ms scan_rows=100000(2639/2500/2406)]
L'opérateur de scan du probe side lit un peu plus de 320 000 lignes depuis le moteur de stockage, bien moins que les 10 millions initiaux. Cela s'explique par le fait que le runtime filter est propagé vers le moteur de stockage, qui utilise les métadonnées d'un lot de données pour filtrer des lots entiers en une seule fois. Cette méthode réduit considérablement la charge d'E/S. Ce type de filtre est particulièrement efficace lorsque la clé de jointure est numérique et que la plage de valeurs du build side est plus étroite que celle du probe side.
Exemple 6 : TopN runtime filter (V4.0+)
Lorsqu'une instruction SQL inclut un opérateur topN, Hologres ne calcule pas tous les résultats. Il génère plutôt un filtre dynamique pour éliminer les données en amont.
SELECT o_orderkey FROM orders ORDER BY o_orderdate LIMIT 5;
Plan d'exécution :
QUERY PLAN
Limit (cost=0.00..116554.70 rows=0 width=8)
-> Sort (cost=0.00..116554.70 rows=100 width=12)
Sort Key: o_orderdate
[id=6 dop=1 time=317/317/317ms rows=5(5/5/5) mem=1/1/1KB open=317/317/317ms get_next=0/0/0ms]
-> Gather (cost=0.00..116554.25 rows=100 width=12)
[20:1 id=100002 dop=1 time=317/317/317ms rows=100(100/100/100) mem=6/6/6KB open=0/0/0ms get_next=317/317/317ms * ]
-> Limit (cost=0.00..116554.25 rows=0 width=12)
-> Sort (cost=0.00..116554.25 rows=150000000 width=12)
Sort Key: o_orderdate
Runtime Filter Sort Column: o_orderdate
[id=3 dop=20 time=318/282/258ms rows=100(5/5/5) mem=96/96/96KB open=318/282/258ms get_next=1/0/0ms]
-> Local Gather (cost=0.00..9.59 rows=150000000 width=12)
[id=2 dop=20 time=316/280/256ms rows=1372205(68691/68610/68498) mem=0/0/0B open=0/0/0ms get_next=316/280/256ms local_dop=1/1/1 * ]
-> Seq Scan on orders (cost=0.00..8.24 rows=150000000 width=12)
Runtime Filter Target Expr: o_orderdate
[id=1 split_count=20 time=286/249/222ms rows=1372205(68691/68610/68498) mem=179/179/179KB open=0/0/0ms get_next=286/249/222ms physical_reads=27074(1426/1353/1294) scan_rows=144867963(7324934/7243398/7172304)]
Query id:[1001003033996040311]
QE version: 2.0
Query Queue: init_warehouse.default_queue
======================cost======================
Total cost:[343] ms
Optimizer cost:[13] ms
Build execution plan cost:[0] ms
Init execution plan cost:[6] ms
Start query cost:[6] ms
- Queue cost: [0] ms
- Wait schema cost:[0] ms
- Lock query cost:[0] ms
- Create dataset reader cost:[0] ms
- Create split reader cost:[0] ms
Get result cost:[318] ms
- Get the first block cost:[318] ms
====================resource====================
Memory: total 7 MB. Worker stats: max 3 MB, avg 3 MB, min 3 MB, max memory worker id: 189*****.
CPU time: total 5167 ms. Worker stats: max 2610 ms, avg 2583 ms, min 2557 ms, max CPU time worker id: 189*****.
DAG CPU time stats: max 5165 ms, avg 2582 ms, min 0 ms, cnt 2, max CPU time dag id: 1.
Fragment CPU time stats: max 5137 ms, avg 1721 ms, min 0 ms, cnt 3, max CPU time fragment id: 2.
Ec wait time: total 90 ms. Worker stats: max 46 ms, max(max) 2 ms, avg 45 ms, min 44 ms, max ec wait time worker id: 189*****, max(max) ec wait time worker id: 189*****.
Physical read bytes: total 799 MB. Worker stats: max 400 MB, avg 399 MB, min 399 MB, max physical read bytes worker id: 189*****.
Read bytes: total 898 MB. Worker stats: max 450 MB, avg 449 MB, min 448 MB, max read bytes worker id: 189*****.
DAG instance count: total 3. Worker stats: max 2, avg 1, min 1, max DAG instance count worker id: 189*****.
Fragment instance count: total 41. Worker stats: max 21, avg 20, min 20, max fragment instance count worker id: 189*****.
Sans le TopN runtime filter, le ScanNode lirait chaque bloc de données de la table orders et le transmettrait au nœud TopN, qui utiliserait ensuite un tri par tas pour conserver les 5 premières lignes rencontrées.
Par exemple, chaque bloc de données contient environ 8 192 lignes. Après le traitement du premier bloc, TopN connaît la valeur o_orderdate classée en 5e position dans ce bloc. Supposons qu'il s'agisse de 1995-01-01. Lorsque le nœud Scan lit le deuxième bloc, il utilise 1995-01-01 comme condition de filtre et n'envoie à TopN que les lignes où o_orderdate <= 1995-01-01. Le seuil est mis à jour dynamiquement. Si la valeur o_orderdate classée en 5e position dans le deuxième bloc est inférieure, TopN remplace l'ancien seuil par la nouvelle valeur.
Consultez la sortie EXPLAIN pour voir le TopN runtime filter généré par l'optimiseur :
-> Limit (cost=0.00..116554.25 rows=0 width=12)
-> Sort (cost=0.00..116554.25 rows=150000000 width=12)
Sort Key: o_orderdate
Runtime Filter Sort Column: o_orderdate
[id=3 dop=20 time=318/282/258ms rows=100(5/5/5) mem=96/96/96KB open=318/282/258ms get_next=1/0/0ms]
La présence de Runtime Filter Sort Column sur le nœud TopN indique que ce nœud génère un TopN runtime filter.
Exemple 7 : Cross join avec ScalarFilter (V4.2+)
SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS t1, t2;
BEGIN;
CREATE TABLE t1 (a int, b int);
CREATE TABLE t2 (a int, b int);
END;
INSERT INTO t1 SELECT t, t FROM generate_series(1, 1000000) t;
INSERT INTO t2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE t1;
ANALYZE t2;
EXPLAIN ANALYZE
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1)
AND a <= (SELECT max(a) FROM t2 WHERE b BETWEEN 0 AND 1);
Points clés du plan d'exécution :
Le nœud Cross Join affiche
Runtime Filter Build Expr: (min(t2.a)), (max(t2.a)).Le ScanNode sur t1 affiche
Runtime Filter Target Expr: (t1.a >= ${1}) AND (t1.a <= ${2}).Pour t1,
scan_rowsreflète le nombre de lignes lues depuis le stockage, tandis querowschute considérablement après l'application du ScalarFilter. Cette différence confirme l'effet du filtre.
Désactivation de la fonctionnalité à titre de comparaison :
SET hg_experimental_generate_runtime_scalar_filter = off;
EXPLAIN ANALYZE
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1)
AND a <= (SELECT max(a) FROM t2 WHERE b BETWEEN 0 AND 1);
Une fois la fonctionnalité désactivée, le plan d'exécution n'affiche plus les champs Runtime Filter Build Expr et Target Expr, et le scan sur t1 redevient un scan complet de la table.
Historique des versions
|
Version |
Nouvelles fonctionnalités |
|
V2.0 |
Prise en charge des runtime filters dans les Hash Join (type Local, incluant les filtres Bloom, In et MinMax). |
|
V2.1 |
Prise en charge des runtime filters avec des conditions de jointure multi-colonnes. |
|
V2.2 |
Prise en charge des runtime filters Global (pour les shuffle joins) ; possibilité de combiner les In filters avec des index bitmap. |
|
V4.0 |
Prise en charge du TopN runtime filter. |
|
V4.2 |
Prise en charge des runtime filters pour les Cross Join (ScalarFilter) afin d'optimiser les scénarios utilisant une sous-requête scalaire pour filtrer une grande table. |