Tous les produits
Search
Centre de documentation

MaxCompute:Clustering par plage

Dernière mise à jour :Sep 02, 2026

Le clustering par plage est une nouvelle méthode de distribution des données selon un ordre de tri global. Cette approche évite les problèmes de déséquilibre des données (data skew) susceptibles de survenir avec le clustering par hachage. Le clustering par plage permet également de créer deux niveaux d'index. Il convient aux scénarios tels que les requêtes par plage basées sur les clés de clustering et les 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 :

  • Pour interroger des données en fonction d'une valeur de colonne spécifique, l'algorithme de hachage localise directement le compartiment de hachage (hash bucket). Ce processus, appelé élagage des compartiments (bucket pruning), réduit la quantité de données analysées. 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, ce qui améliore l'efficacité des requêtes.

  • Lors de la jointure de deux tables via une colonne spécifique dans chaque table, si la colonne de l'une des tables est hachée, l'étape de redistribution (shuffle) peut être supprimée, permettant ainsi 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 les limites suivantes :

  • La création de compartiments via l'algorithme de hachage peut entraîner des problèmes de 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. Une répartition inégale des données d'entrée entre les compartiments provoque des variations considérables de volume de données. Avec le clustering par hachage, chaque compartiment constitue généralement une unité de traitement parallèle. Des différences de volume entre les compartiments sont susceptibles de générer des effets de longue traîne (long tails).

  • L'élagage des compartiments prend uniquement en charge les requêtes d'égalité. Pour les conditions d'inégalité (par exemple, valeurs d'une colonne supérieures à 0), il est impossible de localiser les compartiments contenant les données recherchées. Vous devez alors interroger les données dans tous les compartiments.

  • Pour les requêtes basées sur plusieurs clés de clustering, les performances ne s'améliorent que si toutes les clés de clustering sont disponibles et si toutes les conditions de requête sont des conditions d'égalité.

    Par exemple, pour interroger des données à partir de la table créée avec l'instruction suivante, les performances ne s'améliorent que si la condition de requête est C1=x AND C2=y. Si vous utilisez C1=x ou C2=y comme condition, la requête n'est pas accélérée par le clustering par hachage. 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 recherché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 remédier à ces limitations, MaxCompute propose une nouvelle méthode de clustering des données : le clustering par plage.

Description de la fonctionnalité

Le clustering par plage divise les données en plusieurs plages disjointes basées sur le 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 dupliquées 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 }

Remarque
  • 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 à uniformiser la taille des compartiments, et non celle des plages. Lors du traitement des données, chaque compartiment constitue une unité de traitement parallèle. Une taille de compartiment uniforme évite les problèmes de longue traîne. Cependant, la distribution des données dans chaque plage peut varier. Par conséquent, des tailles de compartiment cohérentes ne signifient pas nécessairement des tailles de plage cohérentes.

MaxCompute met automatiquement en œuvre le processus de clustering par plage. Vous n'avez pas besoin de spécifier manuellement chaque plage, ce qui serait inefficace et irréalisable dans les scénarios de big data. 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. Ainsi, MaxCompute atteint les performances optimales du clustering par plage.

Lors de la création d'une table, spécifiez à la fois RANGE CLUSTERED BY et SORTED BY pour garantir un tri global des données. Ensuite, MaxCompute crée 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.range clustering

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 exclut le Compartiment 2 et le Compartiment 3 sur la base de l'index global, puis interroge les données dans le Compartiment 0 et le Compartiment 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 et c3. Ainsi, vous pouvez interroger les données de la table selon des conditions complexes, telles que c1 = 100 AND c2 > 0 ou c1 = 100 AND c2 = 50 AND c3 < 5. Ce type de requête ne peut pas être réalisé avec le clustering par hachage.

    Important

    Pour une requête basée sur plusieurs clés, les clés de la condition de requête doivent être triées dans l'ordre et seule la dernière clé peut être utilisée pour spécifier une plage de valeurs.

  • Mise en œuvre efficace du tri global.

    Avant l'utilisation du clustering par plage, MaxCompute ne pouvait utiliser qu'une seule instance pour trier globalement les données, ce qui entraînait une faible efficacité. Avec le clustering par plage, les données de chaque plage sont triées simultanément puis combinées, ce qui améliore considérablement l'efficacité.

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é de plage 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. Dans cette instruction, vous devez spécifier le paramètre RANGE CLUSTERED BY. Les paramètres INTO number_of_buckets BUCKETS et SORTED BY sont facultatifs. Nous vous recommandons généralement de spécifier les mêmes valeurs dans SORTED BY et RANGE CLUSTERED BY pour obtenir l'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 spécifié 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 spécifié. Pour éviter les problèmes de déséquilibre des 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 avec de grandes plages de valeurs et peu de valeurs de clés dupliquées. Pour optimiser les performances des requêtes, nous vous conseillons de spécifier dans RANGE CLUSTERED BY les clés d'agrégation ou de filtrage couramment utilisées.

    • SORTED BY

      Spécifie comment trier les champs dans un compartiment. Pour obtenir des performances de requête optimales, nous vous recommandons de spécifier les mêmes valeurs dans SORTED BY et RANGE CLUSTERED BY. Après avoir spécifié SORTED BY, 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 BUCKETS est facultatif pour le clustering par plage. Si vous ne spécifiez pas ce paramètre, MaxCompute détermine automatiquement le nombre de compartiments en fonction du volume de données. Nous vous recommandons généralement de spécifier le nombre de compartiments en fonction de votre situation réelle.

      Comme pour le clustering par hachage, nous vous recommandons de définir le nombre de compartiments en fonction de la taille du compartiment (512 Mo à 1 Go) pour le clustering par plage. Pour les tables extrêmement volumineuses, un grand nombre de compartiments est requis. Toutefois, nous recommandons que le nombre de compartiments dans une table ne dépasse pas 4 000.

Modifier les propriétés de clustering par hachage d'une table

Pour les tables partitionnées, exécutez 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 ne peut modifier que 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 à la table.

    • L'instruction ALTER TABLE prend effet uniquement pour les nouvelles partitions d'une table, y compris les nouvelles partitions générées à l'aide de l'instruction INSERT OVERWRITE. Les nouvelles partitions sont stockées selon les propriétés de clustering. Les formats de stockage des partitions existantes restent inchangés.

    • L'instruction ALTER TABLE prend effet uniquement pour les nouvelles partitions d'une table. Par conséquent, vous ne pouvez pas spécifier de partitions dans cette instruction.

L'instruction ALTER TABLE convient aux tables existantes. Après l'ajout des propriétés de clustering par plage, les nouvelles partitions sont stockées selon 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 consulter les propriétés de la table. Les propriétés de clustering par plage sont affichées dans Extended Info du résultat renvoyé.

DESC EXTENDED <table_name>;

Exemple de sortie :


Owner: ALIYUNxxx                | Project: xxx
  TableComment:

+----------------------------+
| CreateTime:               2018-01-15 16:05:30  |
| LastDDLTime:              2018-01-15 16:05:30  |
| LastModifiedTime:         2018-01-15 16:05:30  |
+----------------------------+
| InternalTable: YES   | Size: 0               |
+----------------------------+
| Native Columns:                                |
+----------------------------+
| Field            | Type     | Label | Comment  |
+----------------------------+
| l_orderkey       | bigint   |       |          |
| l_partkey        | bigint   |       |          |
| l_suppkey        | bigint   |       |          |
| l_linenumber     | bigint   |       |          |
| l_quantity       | double   |       |          |
| l_extendedprice  | double   |       |          |
| l_discount       | double   |       |          |
| l_tax            | double   |       |          |
| l_returnflag     | string   |       |          |
| l_linestatus     | string   |       |          |
| l_shipdate       | string   |       |          |
| l_commitdate     | string   |       |          |
| l_receiptdate    | string   |       |          |
| l_shipinstruct   | string   |       |          |
| l_shipmode       | string   |       |          |
| l_comment        | string   |       |          |
+----------------------------+
| Extended Info:                                 |
+----------------------------+
| TableID:                   xxx                 |
| IsArchived:                false               |
| PhysicalSize:              0                   |
| FileNum:                   0                   |
| ClusterType:               range               |
| BucketNum:                 1024                |
| ClusterColumns:            [l_orderkey]        |
| SortColumns:               [l_orderkey ASC]    |
+----------------------------+

Dans la section Extended Info, ClusterType est range, BucketNum est 1024, ClusterColumns est [l_orderkey] et SortColumns est [l_orderkey ASC], ce qui confirme que la table a été configurée avec succès avec les attributs de clustering par plage.

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>);

Dans la sortie, les quatre attributs ClusterType, BucketNum, ClusterColumns et SortColumns indiquent la configuration de clustering par plage de la partition.


odps@ tpch_100g>desc extended xndai_test_range partition(pt="20180115");

| PartitionSize: 0

| CreateTime:               2018-01-15 16:31:10
| LastDDLTime:              2018-01-15 16:31:10
| LastModifiedTime:         2018-01-15 16:31:10

| IsExstore:          false
| IsArchived:         false
| PhysicalSize:       0
| FileNum:            0
| ClusterType:        range
| BucketNum:          1024
| ClusterColumns:     [c1]
| SortColumns:        [c1 ASC]

Scénarios

Optimisation des requêtes par filtrage

Si le clustering par plage est activé pour une table, les données de la table sont triées globalement. MaxCompute crée automatiquement un index global et des index de fichier basés sur les données triées. Cela améliore l'efficacité du filtrage des données en tirant parti des caractéristiques de stockage des données. 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 id < 3, le système extrait la condition de l'optimiseur et la convertit en la plage de valeurs (-∞, 3). Dans ce cas, le système utilise l'index global pour l'élagage des compartiments afin d'exclure le Compartiment 2 et le Compartiment 3 dont les données ne se trouvent pas dans la plage de valeurs précédente. Ensuite, le système utilise l'index de chaque fichier dans le Compartiment 0 et le Compartiment 1 pour localiser rapidement les données. Ce processus est appelé poussée de prédicat (predicate pushdown), comme illustré dans la figure suivante.过滤查询优化 L'instruction suivante est une instruction TPC-H Query 6 utilisée pour interroger des données à partir d'un ensemble de données de 100 Go après application du clustering par plage. Dans une instruction TPC-H Query 6, une opération d'agrégation est effectuée basée sur un filtrage par plage. Le clustering par plage utilise deux niveaux d'index pour localiser rapidement les données. Ainsi, la durée d'exécution de la requête ainsi que les ressources CPU et mémoire consommées sont considérablement réduites.

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 convertit la table mf_tab en une table avec clustering par plage. Cela vous aide à 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;

Après la conversion de la table en table avec clustering par plage, 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";

Vous pouvez également 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 tirent pleinement parti de la fonctionnalité de tri global du clustering par plage et effectuent une poussée de prédicat pour réduire le nombre d'opérations d'E/S lors de l'analyse de la table, tout en économisant les ressources CPU et mémoire consommées pour le filtrage et le calcul des données.

Si 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 les requêtes de données, k0, k1, ..., km-1 doivent tous être spécifiés dans les conditions et toutes les conditions doivent être des conditions d'égalité. Ainsi, les performances optimales d'accélération des requêtes basées sur les index peuvent être atteintes.

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 peut être réalisée.

  • Si la condition de requête est k1 = 10 AND k2 = 20, l'accélération des requêtes basée sur les index peut être réalisée.

  • Si la condition de requête est k1 = 10 AND k2 < 0, l'accélération des requêtes basée sur les index peut être réalisée.

  • Si la condition de requête est k2 < 0, l'accélération des requêtes basée sur les index ne peut pas être réalisée. En effet, k1 n'est pas spécifié dans la condition de requête.

  • Si la condition de requête est k1 < 0 AND k2 > 0, l'accélération basée sur les index peut être utilisée pour obtenir les données répondant à la condition k1 < 0. Pour les données répondant à la condition k2 > 0, une analyse complète de la table est requise.

Optimisation GROUP BY

Si le clustering par plage est activé pour une table, les données de la table 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 peut être utilisée pour 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 des données à partir de la table, vous pouvez effectuer l'opération GROUP BY sur les données de la table dans une étape de 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;
Remarque

Pour atteindre des performances optimales avec GROUP BY, vous devez spécifier les mêmes clés dans GROUP BY et RANGE CLUSTERED BY.

Optimisation de l'agrégation

L'instruction suivante montre 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 précédentes 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 la colonne a et la colonne b, démarrer trois instances pour agréger les données dans chaque compartiment, puis renvoyer le résultat de sortie.

Cependant, si les valeurs de la colonne a sont distribuées sur plusieurs compartiments, un résultat invalide 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. 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 agréger les tuples. Ainsi, les compartiments sont recréés, comme illustré dans la figure suivante. L'espace entre deux lignes pointillées rouges spécifie la plage de données qui peut ê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 dans la table. Les valeurs obtenues sont enregistrées dans un histogramme. L'histogramme de chaque compartiment est stocké dans les fichiers de métadonnées de cluster. Ce type d'histogramme est appelé histogramme équi-profond (equi-depth histogram).

Remarque

Les tuples sont échantillonnés uniquement si les clés de clustering et les clés de tri sont identiques.

Après obtention de l'histogramme de chaque compartiment, vous pouvez recréer des compartiments pour chaque worker selon les règles suivantes :

  • Les tuples ayant 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.

Sur la base des valeurs limites inférieures de chaque nouveau compartiment, chaque worker peut lire les données dans une plage valide et renvoyer un résultat valide.

Le contenu suivant utilise la table partsupp avec 1 To de données dans un ensemble de données TPC-H pour tester les améliorations de performance. 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;

    Ce qui suit montre le journal d'exécution de la tâche lorsque l'optimisation est désactivée.

    
    resource cost: cpu 9.17 Core * Min, memory 15.15 GB * Min
    inputs:
        yuan_tpch_range_1t.partsupp: 800000000 (3376541880 bytes)
    outputs:
    Job run time: 14.000
    Job run mode: fuxi job
    Job run engine: execution engine
    M1:
        instance count: 128
        run time: 7.000
        instance time:
            min: 2.000, max: 4.000, avg: 2.000
        input records:
            TableScan1: 800000000 (min: 5496040, max: 6780663, avg: 6251772)
        output records:
            StreamLineWrite1: 200001977 (min: 1374023, max: 1695182, avg: 1562958)
        writer dumps:
            StreamLineWrite1: (min: 0, max: 0, avg: 0)
    R2_1:
        instance count: 43
        run time: 14.000
        instance time:
            min: 3.000, max: 4.000, avg: 3.000
        input records:
            StreamLineRead1: 200001977 (min: 4647288, max: 4654888, avg: 4651214)
        output records:
            AdhocSink1: 200000000 (min: 4647242, max: 4654834, avg: 4651168)
        reader dumps:
            StreamLineRead1: (min: 0, max: 0, avg: 0)
  • Exécutez la commande suivante pour activer l'optimisation :

    set odps.optimizer.enable.range.partial.repartitioning=true;

    Exemple de sortie :

    
    resource cost: cpu 4.38 Core * Min, memory 4.38 GB * Min
    inputs:
        yuan_tpch_range_1t.partsupp: 800000000 (18493328320 bytes)
    outputs:
    Job run time: 6.000
    Job run mode: fuxi job
    Job run engine: execution engine
    M1:
        instance count: 128
        run time: 6.000
        instance time:
            min: 1.000, max: 3.000, avg: 2.000
        input records:
            TableScan1: 800000000 (min: 5625876, max: 6259956, avg: 6254874)
        output records:
            AdhocSink1: 200000000 (min: 1406469, max: 1564989, avg: 1563718)

Le résultat du test indique que la vitesse de requête est améliorée de 57 %, l'utilisation du CPU est diminuée de 52 % et l'utilisation de la mémoire est diminuée de 71 % après activation de l'optimisation. Les améliorations de performance varient en fonction du volume de données et des types de requête.

Optimisation des jointures de 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;

    Ensuite, insérez 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 pour les deux tables, les données dans les compartiments des tables peuvent être jointes. Cependant, cette règle ne s'applique pas aux tables avec clustering par plage. Pour deux tables avec clustering par plage devant être jointes, les données dans les compartiments des deux tables ne peuvent pas être directement jointes sur la base des ID de compartiment, même si le nombre de compartiments est identique pour les tables. En effet, la limite de chaque compartiment dans une table avec clustering par plage peut être différente. 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.hash clustering Pour optimiser les jointures entre deux tables avec clustering par plage, recréez des compartiments pour les deux tables en alignant les limites des tables. Ainsi, les limites des données pouvant être lues par chaque instance sont redéfinies.

  • 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.bucket boundary 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 la limite d'une table ayant plus de compartiments avec la limite d'une autre table, et obtient une nouvelle limite pour chaque table, comme illustré dans la figure suivante.bucket优化后 Ainsi, un plan d'exécution sans l'é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 de la colonne a et de la colonne b, puis aligne et recrée des compartiments pour obtenir des limites sur la base desquelles les données sont lues à partir de chaque compartiment de chaque table. Ainsi, un plan d'exécution sans l'étape de redistribution (shuffle) est généré, comme illustré dans la figure précédente.

  • Test de performance

    • Transformation de table

      Utilisez une instruction TPC-H Query 2 pour effectuer des tests sur deux tables nommées PART et PARTSUPP. Chacune des tables contient 1 To de données. Transformez les deux tables en tables avec clustering par plage, et laissez 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 du journal d'exécution de la tâche lorsque l'optimisation est désactivée :

        
        resource cost: cpu 61.64 Core * Min, memory 41.62 GB * Min
        inputs:
            yuan_tpch_range_1t.nation: 25 (1848 bytes)
            yuan_tpch_range_1t.partsupp: 800000000 (7392850104 bytes)
            yuan_tpch_range_1t.region: 5 (1040 bytes)
            yuan_tpch_range_1t.part: 200000000 (1427093008 bytes)
            yuan_tpch_range_1t.supplier: 10000000 (483851352 bytes)
        outputs:
            Job run time: 56.000
            Job run mode: fuxi job
            Job run engine: execution engine
            J11_13:
                instance count: 299
                run time: 52.000
                instance time:
                    min: 1.000, max: 2.000, avg: 1.000
                input records:
                    StreamLineRead13: 637969 (min: 1978, max: 2313, avg: 2133)
                    StreamLineRead7: 159971440 (min: 532818, max: 537639, avg: 535020)
                output records:
                    StreamLineWrite14: 470727 (min: 1473, max: 1676, avg: 1574)
      • Exécutez la commande suivante pour activer l'optimisation :

        set odps.optimizer.enable.range.partial.repartitioning=true;

        Exemple de sortie :

        
        resource cost: cpu 39.81 Core * Min, memory 18.89 GB * Min
        inputs:
            yuan_tpch_range_1t.nation: 25 (1848 bytes)
            yuan_tpch_range_1t.region: 5 (1040 bytes)
            yuan_tpch_range_1t.part: 200000000 (7544753176 bytes)
            yuan_tpch_range_1t.partsupp: 800000000 (22722759616 bytes)
            yuan_tpch_range_1t.supplier: 10000000 (483851352 bytes)
        outputs:
        Job run time: 44.000
        Job run mode: fuxi job
        Job run engine: execution engine
        J11_13:
            instance count: 135
            run time: 40.000
            instance time:
                min: 1.000, max: 3.000, avg: 1.000
            input records:
                StreamLineRead11: 637969 (min: 4516, max: 4962, avg: 4725)
                StreamLineRead7: 159971440 (min: 1181336, max: 1189183, avg: 1184969)
            output records:
                StreamLineWrite12: 470727 (min: 3368, max: 3647, avg: 3486)
            writer dumps:
                StreamLineWrite12: (min: 0, max: 0, avg: 0)
            reader dumps:
                StreamLineRead11: (min: 0, max: 0, avg: 0)
                StreamLineRead7: (min: 0, max: 0, avg: 0)

      Après l'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 diminue d'environ 54,6 %.

Accélération du tri global

Le clustering par plage peut également être utilisé pour l'accélération du 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 un tri global. Cependant, le traitement parallèle ne peut pas être pleinement exploité dans ces scénarios. Vous pouvez utiliser l'étape de partitionnement du clustering par plage pour implémenter un tri global parallèle. Pour le tri global, vous devez échantillonner les données et les diviser en plages, trier les données dans chaque plage en parallèle, puis obtenir le résultat du tri global.

Une fois le tri global terminé, plusieurs compartiments sont toujours inclus dans une table lorsque vous modifiez les propriétés de clustering de la table ou d'une partition de la table. Lors de la consommation des données, les données dans les fichiers doivent être lues en fonction des ID de compartiment pour garantir un tri global.

Par défaut, l'accélération du tri global est désactivée pour les tables avec clustering par plage. Pour activer l'accélération du tri global, 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 :

  • Les coûts de génération des données pour le clustering par plage sont plus élevés que ceux du clustering par hachage. Le clustering par hachage n'est qu'une opération simple de hachage et de tri des données. Cependant, pour le clustering par plage, l'échantillonnage des données, le tri et la combinaison d'histogrammes sont requis. La consommation globale, y compris les durées d'exécution, les coûts CPU et les coûts mémoire, est supérieure à la consommation globale du clustering par hachage. Par conséquent, si le clustering par hachage peut résoudre les problèmes, vous n'avez pas besoin d'utiliser le clustering par plage.

  • Le clustering par plage n'est pas pris en charge dans DYNAMIC PARTITION ou 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. Le clustering par plage 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 peut être réalisée sur la table. Cependant, si range clustered by(a,b) sorted by (b,a) est spécifié pour une table nommée bar, l'optimisation décrite dans cette rubrique ne peut pas être réalisée sur la table.

  • Les clés spécifiées dans JOIN ou GROUP BY doivent être les préfixes ou l'ensemble 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 dans cette rubrique peut être réalisée sur la table uniquement si a, a,b ou a,b,c est spécifié comme clé dans JOIN ou GROUP BY. L'optimisation décrite dans cette rubrique ne peut pas être réalisée si b ou a,c est spécifié comme clé dans JOIN ou GROUP BY.

  • Pour une table partitionnée avec clustering par plage, si vous souhaitez lire des données à partir de deux partitions ou plus de la table, l'optimisation décrite dans cette rubrique ne peut pas être réalisée. L'optimisation décrite dans cette rubrique ne peut être réalisée que sur les tables partitionnées ayant une seule partition et les tables non partitionnées.