Le clustering par plage est une nouvelle méthode de clustering des données qui distribue celles-ci selon un ordre de tri global. Cette approche permet d'éviter les problèmes de déséquilibre des données (data skew) souvent rencontrés avec le clustering par hachage. Elle offre également la possibilité de créer des index à deux niveaux. Le clustering par plage est particulièrement adapté aux scénarios impliquant des requêtes par plage basées sur des clés de clustering ou des requêtes multi-clés. Cette rubrique explique comment utiliser le clustering par plage dans MaxCompute.
Informations générales
Les tables avec clustering par hachage présentent les avantages suivants :
Lorsque vous souhaitez interroger des données en fonction d'une valeur de colonne spécifique, l'algorithme de hachage permet de localiser directement le compartiment (bucket) concerné. Ce processus est appelé élagage des compartiments (bucket pruning). Si les données du compartiment sont stockées de manière triée, vous pouvez utiliser des index pour localiser plus précisément les données. Cela réduit la quantité de données analysées et améliore l'efficacité des requêtes.
Si vous effectuez une jointure entre deux tables sur une colonne spécifique de chacune, et que la colonne de l'une des tables est hachée, l'étape de redistribution (shuffle) peut être supprimée. Cela permet d'économiser des ressources de calcul.
Pour plus d'informations sur la fonctionnalité de clustering par hachage, consultez Clustering par hachage.
Le clustering par hachage présente toutefois les limites suivantes :
L'utilisation de l'algorithme de hachage pour créer des compartiments peut entraîner un déséquilibre des données. Semblables aux problèmes de déséquilibre lors des jointures, ces déséquilibres sont inhérents à l'algorithme de hachage. Si les données d'entrée ne sont pas réparties uniformément entre les compartiments, la quantité de données peut varier considérablement d'un compartiment à l'autre. Avec le clustering par hachage, chaque compartiment constitue généralement une unité de traitement parallèle. Des différences importantes de volume de données entre les compartiments peuvent alors provoquer des effets de longue traîne (long tails).
L'élagage des compartiments ne prend en charge que les requêtes d'égalité. Pour les requêtes basées sur des conditions d'inégalité (par exemple, lorsque les valeurs d'une colonne sont supérieures à 0), il est impossible de déterminer quels compartiments contiennent les données recherchées. Dans ce cas, vous devez analyser l'ensemble des compartiments.
-
Pour les requêtes utilisant plusieurs clés de clustering, les performances ne sont améliorées que si toutes les clés de clustering sont disponibles et si toutes les conditions de requête sont des égalités.
Par exemple, supposons que vous souhaitiez interroger des données depuis une table créée à l'aide de l'instruction suivante. L'amélioration des performances n'est effective que si la condition de requête est
C1=x AND C2=y. Si vous utilisezC1=xouC2=ycomme condition, l'accélération basée sur le clustering par hachage ne s'applique pas. En effet, les valeurs de hachage des clés sont combinées par paires lors des requêtes. Sans cette combinaison, il est impossible de localiser le compartiment contenant les données ou d'effectuer l'élagage des compartiments.CREATE TABLE T2 (C1 int, C2 int, C3 string) CLUSTERED BY (C1, C2) SORTED by (C1, C2) INTO 1024 BUCKETS;
Pour pallier ces limitations, MaxCompute propose une nouvelle méthode de clustering des données appelée clustering par plage.
Description de la fonctionnalité
Le clustering par plage divise les données en plusieurs plages disjointes en se basant sur un tri complet des clés de clustering. Chaque plage est considérée comme un compartiment et doit respecter les deux conditions suivantes :
Les valeurs en double sont stockées dans le même compartiment.
Le nombre de valeurs dans chaque compartiment est approximativement équivalent.
L'instruction suivante crée une table nommée T.
CREATE TABLE T (C1 int)
RANGE CLUSTERED BY (C1)
SORTED BY (c1)
INTO 3 BUCKETS;
Les valeurs de la colonne C1 sont { 1, 8, -3, 2, 4, 1, 1, 3, 8, 20, -8, 9 }.
Après activation du clustering par plage, les compartiments suivants sont obtenus :
Compartiment 0 : { -8, -3, 1, 1, 1 }
Compartiment 1 : { 2, 3, 4 }
Compartiment 2 : { 8, 8, 9, 20 }
Les plages représentées par les compartiments peuvent être disjointes. Par exemple, la plage du compartiment 1 est
[2, 4]et celle du compartiment 2 est[8, 20]. Aucune valeur ne se trouve dans la plage(4, 8).Le clustering par plage vise à obtenir des tailles de compartiments similaires plutôt que des tailles de plages similaires. Lors du traitement des données, chaque compartiment sert d'unité pour le traitement parallèle. Des tailles de compartiments uniformes évitent les problèmes de longue traîne. Cependant, la distribution des données au sein de chaque plage peut varier. Ainsi, des tailles de compartiments constantes n'impliquent pas nécessairement des tailles de plages constantes.
Le processus de clustering par plage est automatiquement géré par MaxCompute. Vous n'avez pas besoin de spécifier manuellement chaque plage. Dans les scénarios de big data, la configuration manuelle des plages n'est ni efficace ni réalisable. MaxCompute trie et échantillonne automatiquement les données, crée un histogramme basé sur la distribution des données de chaque plage, puis combine et calcule l'histogramme de chaque plage. Cette approche permet à MaxCompute d'atteindre des performances optimales pour le clustering par plage.
Lors de la création d'une table, vous pouvez spécifier à la fois RANGE CLUSTERED BY et SORTED BY afin de garantir un tri global des données. MaxCompute crée alors automatiquement deux niveaux d'index : un index global et un index de fichier. Ces index permettent de localiser et de rechercher rapidement les valeurs clés, comme illustré dans la figure suivante.
Le clustering par plage présente les avantages suivants par rapport au clustering par hachage :
-
Prise en charge des requêtes par plage.
Par exemple, si la condition de requête est
c < 3, le système peut exclure les compartiments 2 et 3 grâce à l'index global, puis interroger les données uniquement dans les compartiments 0 et 1. Avec le clustering par hachage, l'élagage des compartiments n'est possible que pour les requêtes d'égalité. -
Prise en charge des requêtes multi-clés.
Par exemple, si vous spécifiez
RANGE CLUSTERED BY (c1, c2, c3) SORTED BY (c1, c2, c3)lors de la création d'une table, le clustering par plage et le stockage des données s'effectuent dans l'ordre c1, c2, puis c3. Vous pouvez ainsi interroger les données de la table avec des conditions complexes, telles quec1 = 100 AND c2 > 0ouc1 = 100 AND c2 = 50 AND c3 < 5. Ce type de requête n'est pas réalisable avec le clustering par hachage.ImportantPour une requête basée sur plusieurs clés, les clés présentes dans la condition de requête doivent être ordonnées séquentiellement, et seule la dernière clé peut servir à spécifier une plage de valeurs.
-
Mise en œuvre efficace du tri global.
Avant l'utilisation du clustering par plage, MaxCompute ne pouvait effectuer un tri global des données qu'à l'aide d'une seule instance, ce qui limitait l'efficacité. Avec le clustering par plage, les données de chaque plage peuvent être triées en parallèle puis combinées, améliorant significativement les performances.
Remarques d'utilisation
La syntaxe du clustering par plage est similaire à celle du clustering par hachage. La différence réside dans le fait que le mot-clé range et le nombre de compartiments sont facultatifs pour le clustering par plage.
Créer une table avec clustering par plage
Utilisez l'instruction CREATE TABLE pour créer une table avec clustering par plage. Vous devez spécifier le paramètre RANGE CLUSTERED BY. Les paramètres INTO number_of_buckets BUCKETS et SORTED BY sont facultatifs. Il est généralement recommandé de spécifier les mêmes valeurs pour SORTED BY et RANGE CLUSTERED BY afin d'obtenir un effet d'optimisation optimal.
-
Syntaxe
CREATE TABLE [IF NOT EXISTS] <table_name> [(<col_name> data_type [comment <col_comment>], ...)] [comment table_comment] [PARTITIONED BY (<col_name> data_type [comment <col_comment>], ...)] [RANGE CLUSTERED BY (<col_name> [, <col_name>, ...]) [SORTED BY (<col_name> [ASC | DESC] [, <col_name> [ASC | DESC] ...])] [INTO <number_of_buckets> BUCKETS]] [AS select_statement] -
Exemples
-
Table non partitionnée
CREATE TABLE T1 (a string, b string, c int) RANGE CLUSTERED BY (c) SORTED by (c) INTO 1024 BUCKETS; -
Table partitionnée
CREATE TABLE T1 (a string, b string, c int) PARTITIONED BY (dt int) RANGE CLUSTERED BY (c) SORTED by (c) INTO 1024 BUCKETS;
-
-
Paramètres
-
RANGE CLUSTERED BY
Spécifie les clés pour le clustering par plage. Après avoir défini ce paramètre, MaxCompute trie et échantillonne une ou plusieurs colonnes de données dans des plages appropriées, en fonction du nombre de compartiments indiqué. Pour éviter les déséquilibres de données et les points chauds, et pour améliorer les performances des requêtes parallèles, nous vous recommandons de spécifier dans RANGE CLUSTERED BY des colonnes présentant de larges plages de valeurs et peu de valeurs de clés en double. Pour optimiser les performances des requêtes, privilégiez les clés d'agrégation ou de filtrage couramment utilisées.
-
SORTED BY
Indique comment trier les champs au sein d'un compartiment. Pour obtenir des performances de requête optimales, il est conseillé de spécifier les mêmes valeurs pour SORTED BY et RANGE CLUSTERED BY. Une fois SORTED BY défini, MaxCompute génère automatiquement un index global et des index de fichier, et utilise ces index pour accélérer les requêtes.
-
INTO number_of_buckets BUCKETS
Contrairement au clustering par hachage,
INTO number_of_buckets BUCKETSest facultatif pour le clustering par plage. Si vous omettez ce paramètre, MaxCompute détermine automatiquement le nombre de compartiments en fonction du volume de données. Il est généralement recommandé de définir le nombre de compartiments en fonction de votre situation réelle.Comme pour le clustering par hachage, nous vous suggérons de définir le nombre de compartiments en fonction de la taille cible d'un compartiment (entre 512 Mo et 1 Go). Pour les tables extrêmement volumineuses, un grand nombre de compartiments sera nécessaire. Toutefois, il est recommandé de ne pas dépasser 4 000 compartiments par table.
-
Modifier les propriétés de clustering par hachage d'une table
Pour les tables partitionnées, vous pouvez exécuter l'instruction ALTER TABLE pour ajouter ou supprimer les propriétés de clustering par plage.
-
Syntaxe
-- Change a table to a range-clustered table. ALTER TABLE <table_name> [RANGE CLUSTERED BY (<col_name> [, <col_name>, ...]) [SORTED BY (<col_name> [ASC | DESC] [, <col_name> [ASC | DESC] ...])] [INTO <number_of_buckets> BUCKETS]; -- Change a range-clustered table to a non-range-clustered table. ALTER TABLE <table_name> NOT CLUSTERED; -
Remarques d'utilisation
L'instruction ALTER TABLE permet uniquement de modifier les propriétés de clustering d'une table partitionnée. Pour une table non partitionnée, les propriétés de clustering ne peuvent pas être modifiées après leur ajout initial.
L'instruction ALTER TABLE s'applique uniquement aux nouvelles partitions d'une table, y compris celles générées via l'instruction INSERT OVERWRITE. Les nouvelles partitions sont stockées selon les propriétés de clustering définies. Le format de stockage des partitions existantes reste inchangé.
L'instruction ALTER TABLE s'appliquant uniquement aux nouvelles partitions, vous ne pouvez pas spécifier de partitions individuelles dans cette instruction.
L'instruction ALTER TABLE convient aux tables existantes. Une fois les propriétés de clustering par plage ajoutées, les nouvelles partitions seront stockées conformément à ces propriétés.
Vérifier explicitement les propriétés de la table
Après avoir créé une table avec clustering par plage, exécutez l'instruction suivante pour afficher les propriétés de la table. Les propriétés de clustering par plage apparaissent dans la section Extended Info du résultat renvoyé.
DESC EXTENDED <table_name>;
La figure suivante illustre un exemple de résultat renvoyé.
Pour une table partitionnée, vous pouvez également exécuter l'instruction suivante pour consulter les propriétés de clustering d'une partition.
DESC EXTENDED <table_name> partition(<pt_spec>);
La figure suivante montre un exemple de résultat renvoyé.
Scénarios
Optimisation des requêtes par filtrage
Si le clustering par plage est activé pour une table, les données de celle-ci sont triées globalement. MaxCompute crée automatiquement un index global et des index de fichier basés sur ces données triées. Cela améliore l'efficacité du filtrage des données en tirant parti des caractéristiques de stockage. Vous pouvez utiliser le clustering par plage pour optimiser les requêtes d'égalité et les requêtes par plage.
Par exemple, pour une condition de requête simple telle que id < 3, le système extrait la condition depuis l'optimiseur et la convertit en plage de valeurs (-∞, 3). Le système peut alors utiliser l'index global pour l'élagage des compartiments, excluant ainsi les compartiments 2 et 3 dont les données ne se trouvent pas dans cette plage. Ensuite, il utilise l'index de chaque fichier dans les compartiments 0 et 1 pour localiser rapidement les données. Ce processus est appelé propagation des prédicats (predicate pushdown), comme illustré dans la figure suivante.
L'instruction suivante correspond à la requête TPC-H Query 6, utilisée pour interroger un jeu de données de 100 Go après application du clustering par plage. Dans cette requête, une opération d'agrégation est effectuée sur la base d'un filtrage par plage. Le clustering par plage exploite les deux niveaux d'index pour localiser rapidement les données, réduisant ainsi considérablement la durée d'exécution de la requête ainsi que les ressources CPU et mémoire consommées.
select sum(l_extendedprice * l_discount) as revenue
from tpch_lineitem l
where l_shipdate >= '1994-01-01'
and l_shipdate < '1995-01-01'
and l_discount >= 0.05
and l_discount <= 0.07
and l_quantity < 24;
Requêtes multi-clés
Dans cet exemple, l'instruction suivante transforme la table mf_tab en une table avec clustering par plage, afin de mieux comprendre les requêtes multi-clés.
ALTER TABLE mf_project.mf_tab
RANGE CLUSTERED BY (project_name, name)
SORTED BY (project_name, name)
INTO 1024 BUCKETS;
Une fois la table transformée, vous pouvez effectuer des requêtes d'agrégation au niveau du projet. Exemple d'instruction :
SELECT COUNT(*)
from mf_project.mf_tab
WHERE project_name="xxxdw"
AND ds="20180115"
AND type="TABLE";
Vous pouvez également utiliser plusieurs clés pour localiser précisément une table. Exemple d'instruction :
SELECT count(*)
from mf_project.mf_tab
WHERE project_name="xxxdw"
AND name="adm_ctu_cle_kba_midun_trade_dd"
AND type="TABLE";
Il est aussi possible d'utiliser plusieurs clés pour des requêtes par plage. L'instruction suivante interroge les tables dont les noms commencent par adm.
SELECT count(*)
from mf_project.mf_tab
WHERE project_name="xxxdw"
AND name>="adm"
AND name < "adn"
AND type="TABLE";
Toutes les requêtes précédentes exploitent pleinement la fonctionnalité de tri global du clustering par plage et effectuent une propagation des prédicats pour réduire le nombre d'opérations d'E/S lors de l'analyse des tables, économisant ainsi les ressources CPU et mémoire utilisées pour le filtrage et le calcul des données.
Lorsque plusieurs clés sont utilisées pour le clustering par plage, certaines exigences doivent être respectées. Pour RANGE CLUSTERED BY k0, k1, ..., kn dans une instruction de création de table, si km est utilisé pour interroger les données, k0, k1, ..., km-1 doivent tous être spécifiés dans les conditions, et toutes ces conditions doivent être des égalités. Cela permet d'obtenir des performances optimales pour l'accélération des requêtes basée sur les index.
Par exemple, k1, k2 sont des clés de clustering dans une table nommée T.
Si la condition de requête est
k1 < 5, l'accélération des requêtes basée sur les index est effective.Si la condition de requête est
k1 = 10 AND k2 = 20, l'accélération des requêtes basée sur les index est effective.Si la condition de requête est
k1 = 10 AND k2 < 0, l'accélération des requêtes basée sur les index est effective.Si la condition de requête est
k2 < 0, l'accélération des requêtes basée sur les index n'est pas possible, car k1 n'est pas spécifiée dans la condition.Si la condition de requête est
k1 < 0 AND k2 > 0, l'accélération basée sur les index permet d'obtenir les données correspondant à la conditionk1 < 0. Pour les données correspondant à la conditionk2 > 0, une analyse complète de la table est nécessaire.
Optimisation de GROUP BY
Si le clustering par plage est activé pour une table, les données sont triées globalement. Les clés ayant les mêmes valeurs sont placées dans le même compartiment lors du clustering par plage. Cette propriété physique des données permet de supprimer l'étape de redistribution (shuffle) lors des opérations d'agrégation.
Par exemple, une table nommée T est créée à l'aide de l'instruction CREATE TABLE suivante. Si vous interrogez les données de cette table, vous pouvez effectuer l'opération GROUP BY sur les données de la table durant l'étape map.
CREATE TABLE T (department int, team string, employee string)
RANGE CLUSTERED BY (department, team)
SORTED BY (c1, c2)
INTO 1024 BUCKETS;
SELECT COUNT(*) from T GROUP BY department, team;
Pour obtenir des performances optimales avec GROUP BY, vous devez spécifier les mêmes clés dans GROUP BY et RANGE CLUSTERED BY.
Optimisation des agrégations
L'instruction suivante présente la structure de données de la table foo.
create table foo(a bigint, b bigint, c bigint)
range clustered by (a,b)
sorted by(a,b) into 3 buckets;
Les données stockées dans les compartiments de la table foo se situent dans les plages suivantes :
Bucket 0: [1,1 : 3,3]
Bucket 1: [5,5 : 7,7]
Bucket 2: [8,8 : 9,9]
Les plages de compartiments ci-dessus sont spécifiées au format Bucket N: [lower bound values : upper bound values]. Si les données sont agrégées par la colonne a, les plages de compartiments suivantes sont utilisées à la place :
Bucket 0: [1 : 3]
Bucket 1: [5 : 7]
Bucket 2: [8 : 9]
Vous pouvez générer directement un plan d'exécution pour les opérations d'agrégation par les colonnes a et b, démarrer trois instances pour agréger les données de chaque compartiment, puis renvoyer le résultat.
Cependant, si les valeurs de la colonne a sont réparties sur plusieurs compartiments, un résultat incorrect est renvoyé. Exemple :
Bucket 0: [1,1 : 3,3]
Bucket 1: [3,5 : 7,7]
Bucket 2: [7,8 : 9,9]
La colonne a possède deux valeurs, 3 et 7, qui sont stockées séparément dans deux compartiments différents. Pour obtenir un résultat valide, vous devez placer les tuples ayant la même valeur dans la colonne a dans la même instance et les agréger. Cela entraîne une recréation des compartiments, comme illustré dans la figure suivante. L'espace entre deux lignes pointillées rouges définit la plage de données pouvant être lue par chaque instance.
Des histogrammes sont requis pour le clustering par plage. Pour une table avec clustering par plage, si les clés de clustering et les clés de tri sont identiques, le worker correspondant à chaque compartiment échantillonne un tuple toutes les 10 000 lignes pour obtenir les valeurs des clés de clustering lors de l'insertion des données. Les valeurs obtenues sont enregistrées dans un histogramme. L'histogramme de chaque compartiment est stocké dans les fichiers de métadonnées du cluster. Ce type d'histogramme est appelé histogramme équi-profond.
L'échantillonnage des tuples n'a lieu que si les clés de clustering et les clés de tri sont identiques.
Une fois l'histogramme de chaque compartiment obtenu, vous pouvez recréer des compartiments pour chaque worker selon les règles suivantes :
Les tuples partageant la même clé de regroupement sont stockés dans le même compartiment.
Les données sont réparties uniformément entre les compartiments.
En se basant sur les valeurs limites inférieures de chaque nouveau compartiment, chaque worker peut lire les données dans une plage valide et renvoyer un résultat correct.
Le contenu suivant utilise la table partsupp contenant 1 To de données dans un jeu de données TPC-H pour tester les améliorations de performances. Exécutez l'instruction suivante pour transformer la table partsupp en une table avec clustering par plage :
CREATE TABLE partsupp ( PS_PARTKEY BIGINT NOT NULL,
PS_SUPPKEY BIGINT NOT NULL,
PS_AVAILQTY BIGINT NOT NULL,
PS_SUPPLYCOST DECIMAL(15,2) NOT NULL,
PS_COMMENT VARCHAR(199) NOT NULL)
RANGE CLUSTERED BY(PS_PARTKEY, PS_SUPPKEY)
SORTED BY(PS_PARTKEY, PS_SUPPKEY) INTO 128 BUCKETS;
Exécutez l'instruction de requête suivante pour effectuer le test :
SELECT ps_partkey, count(*) c FROM partsupp GROUP BY ps_partkey;
-
Exécutez la commande suivante pour désactiver l'optimisation :
set odps.optimizer.enable.range.partial.repartitioning=false;Exemple de sortie :

-
Exécutez la commande suivante pour activer l'optimisation :
set odps.optimizer.enable.range.partial.repartitioning=true;Exemple de sortie :

Le résultat du test indique qu'après activation de l'optimisation, la vitesse de requête augmente de 57 %, l'utilisation du CPU diminue de 52 % et l'utilisation de la mémoire baisse de 71 %. Les gains de performance varient en fonction du volume de données et des types de requêtes.
Optimisation des jointures pour les tables avec clustering par plage
-
Dans cet exemple, créez deux tables à l'aide des instructions suivantes :
create table t1(a bigint, b bigint, c bigint, d bigint) range clustered by(a,b,c) sorted by(a,b,c) into 3 buckets; create table t2(a bigint, b bigint, c bigint, d bigint) range clustered by(a,b,c) sorted by(a,b,c) into 3 buckets;Insérez ensuite des données différentes dans les deux tables.
Pour deux tables avec clustering par hachage devant être jointes, si le nombre de compartiments est identique, les données des compartiments peuvent être jointes directement. Cette règle ne s'applique pas aux tables avec clustering par plage. Pour deux tables avec clustering par plage, les données des compartiments ne peuvent pas être jointes directement sur la base des ID de compartiment, même si le nombre de compartiments est identique. En effet, les limites de chaque compartiment dans une table avec clustering par plage peuvent différer. Si vous joignez deux tables avec clustering par plage, un plan d'exécution incluant l'étape de redistribution (shuffle) est toujours généré, comme illustré dans la figure suivante.
Pour optimiser les jointures entre deux tables avec clustering par plage, recréez des compartiments pour les deux tables en alignant leurs limites. Cela permet de redéfinir les limites des données lisibles par chaque instance. -
Créez deux tables.
create table t1(a bigint, b bigint, c bigint, d bigint) range clustered by(a,b,c) sorted by(a,b,c) into 5 buckets; create table t2(a bigint, b bigint, c bigint, d bigint) range clustered by(a,b,c) sorted by(a,b,c) into 3 buckets;Après insertion d'une certaine quantité de données dans les tables, les limites des compartiments sont définies, comme illustré dans la figure suivante.
Exemple de requête 1 :SELECT * FROM t1 JOIN t2 ON t1.a=t2.a AND t1.b=t2.b AND t1.c=t2.c;L'optimiseur aligne les limites de la table comportant plus de compartiments sur celles de l'autre table, obtenant ainsi de nouvelles limites pour chaque table, comme illustré dans la figure suivante.
De cette manière, un plan d'exécution sans étape de redistribution (shuffle) est généré.
Exemple de requête 2 :SELECT * FROM t1 JOIN t2 ON t1.a=t2.a AND t1.b=t2.b;L'optimiseur crée des compartiments pour chaque table en fonction des colonnes a et b, puis aligne et recrée les compartiments pour obtenir des limites permettant de lire les données de chaque compartiment de chaque table. Un plan d'exécution sans étape de redistribution (shuffle) est ainsi généré, comme illustré dans la figure précédente.
-
Tests de performance
-
Transformation de table
Effectuez un test sur deux tables nommées PART et PARTSUPP à l'aide d'une instruction TPC-H Query 2. Chacune des tables contient 1 To de données. Transformez les deux tables en tables avec clustering par plage, tout en laissant les autres tables inchangées.
CREATE TABLE PARTSUPP ( PS_PARTKEY BIGINT NOT NULL, PS_SUPPKEY BIGINT NOT NULL, PS_AVAILQTY BIGINT NOT NULL, PS_SUPPLYCOST DECIMAL(15,2) NOT NULL, PS_COMMENT VARCHAR(199) NOT NULL) RANGE CLUSTERED BY(PS_PARTKEY, PS_SUPPKEY) SORTED BY(PS_PARTKEY, PS_SUPPKEY) INTO 128 BUCKETS; CREATE TABLE PART ( P_PARTKEY BIGINT NOT NULL, P_NAME VARCHAR(55) NOT NULL, P_MFGR CHAR(25) NOT NULL, P_BRAND CHAR(10) NOT NULL, P_TYPE VARCHAR(25) NOT NULL, P_SIZE BIGINT NOT NULL, P_CONTAINER CHAR(10) NOT NULL, P_RETAILPRICE DECIMAL(15,2) NOT NULL, P_COMMENT VARCHAR(23) NOT NULL) RANGE CLUSTERED BY(P_PARTKEY) SORTED BY(P_PARTKEY) INTO 64 BUCKETS;Interrogez les données à l'aide de l'instruction TPC-H Query 2 suivante :
select s_acctbal, s_name, n_name, p_partkey, p_mfgr, s_address, s_phone, s_comment from part, supplier, partsupp, nation, region where p_partkey = ps_partkey and s_suppkey = ps_suppkey and p_size = 15 and p_type like '%BRASS' and s_nationkey = n_nationkey and n_regionkey = r_regionkey and r_name = 'EUROPE' and ps_supplycost = (select min(ps_supplycost) from partsupp, supplier, nation, region where p_partkey = ps_partkey and s_suppkey = ps_suppkey and s_nationkey = n_nationkey and n_regionkey = r_regionkey and r_name = 'EUROPE') order by s_acctbal desc, n_name, s_name, p_partkey limit 100; -
Résultats des tests
-
Exécutez la commande suivante pour désactiver l'optimisation :
set odps.optimizer.enable.range.partial.repartitioning=false;Exemple de sortie :

-
Exécutez la commande suivante pour activer l'optimisation :
set odps.optimizer.enable.range.partial.repartitioning=true;Exemple de sortie :

Après optimisation, deux étapes sont supprimées, la vitesse de requête s'améliore d'environ 21,4 %, l'utilisation du CPU diminue d'environ 35,4 % et l'utilisation de la mémoire baisse d'environ 54,6 %.
-
-
Accélération du tri global
Le clustering par plage peut également servir à accélérer le tri global. Dans les scénarios courants utilisant ORDER BY, toutes les données triées sont distribuées à la même instance pour garantir le tri global. Toutefois, ces scénarios n'exploitent pas pleinement le traitement parallèle. Vous pouvez utiliser l'étape de partitionnement du clustering par plage pour mettre en œuvre un tri global parallèle. Pour le tri global, vous devez échantillonner les données, les diviser en plages, trier les données de chaque plage en parallèle, puis obtenir le résultat du tri global.
Une fois le tri global terminé, plusieurs compartiments subsistent dans une table lorsque vous modifiez les propriétés de clustering de la table ou d'une partition. Lors de la consommation des données, les fichiers doivent être lus en fonction des ID de compartiment pour garantir le tri global.
Par défaut, l'accélération du tri global est désactivée pour les tables avec clustering par plage. Pour l'activer, exécutez la commande suivante :
set odps.optimizer.distribute.ordering.enable=true;
Limites et remarques d'utilisation
Par rapport au clustering par hachage, le clustering par plage présente les limites suivantes :
Le coût de génération des données pour le clustering par plage est plus élevé que pour le clustering par hachage. Le clustering par hachage consiste simplement en un hachage et un tri des données. En revanche, le clustering par plage nécessite un échantillonnage des données, un tri et une combinaison d'histogrammes. La consommation globale, incluant les durées d'exécution ainsi que les coûts CPU et mémoire, est supérieure à celle du clustering par hachage. Par conséquent, si le clustering par hachage suffit à résoudre vos problèmes, inutile d'utiliser le clustering par plage.
Le clustering par plage n'est pas pris en charge dans DYNAMIC PARTITION ni INSERT INTO.
Le clustering par plage est pris en charge uniquement pour les opérations de jointure suivantes : inner join, left outer join, right outer join et semi join. Il n'est pas pris en charge pour anti-join ou full outer join.
Pour les tables avec clustering par plage, les clés spécifiées dans RANGE CLUSTERED BY doivent être identiques à celles spécifiées dans SORTED BY. Par exemple, si
range clustered by (a,b) sorted by (a,b)est spécifié pour une table nommée foo lors de sa création, l'optimisation décrite dans cette rubrique s'applique à la table. En revanche, sirange clustered by(a,b) sorted by (b,a)est spécifié pour une table nommée bar, cette optimisation ne s'applique pas.Les clés spécifiées dans JOIN ou GROUP BY doivent correspondre aux préfixes ou à l'intégralité des clés spécifiées dans RANGE CLUSTERED BY. Par exemple,
range clustered by(a,b,c) sorted by(a,b,c)est spécifié dans une instruction de création de table. L'optimisation décrite ici s'applique uniquement sia,a,boua,b,cest spécifié comme clé dans JOIN ou GROUP BY. L'optimisation ne s'applique pas siboua,cest spécifié comme clé.Pour une table partitionnée avec clustering par plage, si vous souhaitez lire les données de deux partitions ou plus, l'optimisation décrite ici ne s'applique pas. Elle n'est effective que pour les tables partitionnées ayant une seule partition et pour les tables non partitionnées.