レンジクラスタリングは、データをグローバルにソートされた順序で分散配置する新しいデータクラスタリング方法です。レンジクラスタリングを使用すると、ハッシュクラスタリングによって発生する可能性があるデータスキューの問題を防止できます。また、レンジクラスタリングでは 2 段階のインデックスを作成できます。レンジクラスタリングは、クラスターキーに基づく範囲クエリや複数キー クエリなどのシナリオに適しています。本トピックでは、MaxCompute でレンジクラスタリングを使用する方法について説明します。
背景情報
ハッシュクラスタリングテーブルには、次の利点があります。
-
特定の列値に基づいてデータをクエリする場合、ハッシュアルゴリズムを使用してハッシュバケットを直接特定できます。この処理はバケットプルーニングと呼ばれます。バケット内のデータがソートされた状態で格納されている場合は、さらにインデックスを使用してデータを特定できます。これによりスキャンされるデータ量が削減され、クエリ効率が向上します。
-
2 つのテーブルをそれぞれの特定の列で結合し、かつ一方のテーブルの列がハッシュ化されている場合、シャッフル手順を省略できます。これにより計算リソースを節約できます。
ハッシュクラスタリング機能の詳細については、「ハッシュクラスタリング」をご参照ください。
ハッシュクラスタリングには、次の制限があります。
-
ハッシュアルゴリズムを使用してバケットを作成すると、データスキューの問題が発生する可能性があります。結合スキューの問題と同様に、データスキューの問題はハッシュアルゴリズムに固有のものです。入力データがバケット間で不均一に分布している場合、データスキューの問題が発生する可能性があります。その結果、バケット内のデータ量が大きく異なります。ハッシュクラスタリングを使用する場合、多くの場合、各バケットが並列処理の単位になります。バケットごとのデータ量が異なると、ロングテールが発生しやすくなります。
-
バケットプルーニングでサポートされるのは等価クエリのみです。たとえば、列の値が 0 より大きいといった非等価条件に基づいてデータをクエリする場合、対象データが格納されているバケットを特定できません。この場合、すべてのバケットからデータをクエリする必要があります。
-
複数のクラスターキーに基づくクエリでは、すべてのクラスターキーが利用可能で、かつすべてのクエリ条件が等価条件である場合にのみ、クエリパフォーマンスを向上させることができます。
たとえば、次のステートメントを使用して作成されたテーブルからデータをクエリするとします。クエリ条件として
C1=x AND C2=yを使用した場合にのみ、クエリパフォーマンスを向上させることができます。クエリ条件としてC1=xまたはC2=yを使用した場合、ハッシュクラスタリングに基づいてクエリを高速化することはできません。これは、キーがクエリに使用される際にキーのハッシュ値が結合されるためです。キーのハッシュ値が結合されない場合、クエリ対象のデータが格納されているバケットを特定できず、バケットプルーニングを実行することもできません。CREATE TABLE T2 (C1 INT, C2 INT, C3 STRING) CLUSTERED BY (C1, C2) SORTED BY (C1, C2) INTO 1024 BUCKETS;
これらの制限に対処するため、MaxCompute はレンジクラスタリングと呼ばれる新しいデータクラスタリング方法を提供します。
機能の説明
レンジクラスタリングは、クラスターキーの完全なソートに基づき、データを複数の互いに重ならない範囲に分割します。各範囲はバケットと見なされ、次の両方の条件を満たす必要があります。
-
重複する値は同じバケットに格納されます。
-
各バケット内の値の数はほぼ同等になります。
次のサンプルステートメントは、T という名前のテーブルを作成します。
CREATE TABLE T (C1 INT)
RANGE CLUSTERED BY (C1)
SORTED BY (c1)
INTO 3 BUCKETS;
C1 列の値は { 1, 8, -3, 2, 4, 1, 1, 3, 8, 20, -8, 9 } です。
レンジクラスタリングを有効にすると、次のバケットが得られます。
-
Bucket 0: { -8, -3, 1, 1, 1 }
-
Bucket 1: { 2, 3, 4 }
-
Bucket 2: { 8, 8, 9, 20 }
-
バケットが表す範囲は、不連続である場合があります。たとえば、バケット 1 の範囲は
[2, 4]で、バケット 2 の範囲は[8, 20]です。(4, 8)の範囲には値がありません。 -
レンジクラスタリングは、範囲サイズをほぼ同等にするのではなく、バケットサイズをほぼ同等にすることを目的としています。データを処理する際、各バケットが並列処理の単位になります。バケットサイズを同じにすることで、ロングテールの問題を防止できます。ただし、各範囲内のデータ分布は同一ではない場合があります。そのため、バケットサイズが一定であっても、範囲サイズが一定であることを意味しません。
レンジクラスタリングの処理は MaxCompute によって自動的に実行されます。各範囲を手動で指定する必要はありません。ビッグデータのシナリオでは、範囲の手動設定は効率的でも現実的でもありません。MaxCompute はデータを自動的にソートしてサンプリングし、各範囲のデータ分布に基づいてヒストグラムを作成し、その後、各範囲のヒストグラムを結合して計算します。これにより、MaxCompute はレンジクラスタリングの最適なパフォーマンスを実現します。
テーブルを作成するときに RANGE CLUSTERED BY と SORTED BY の両方を指定してデータをグローバルにソートすると、MaxCompute は自動的にグローバルインデックスとファイルインデックスという 2 つのレベルのインデックスを作成します。これらのインデックスは、次の図に示すように、キー値を迅速に特定および検索するために使用されます。
レンジクラスタリングには、ハッシュクラスタリングと比較して次の利点があります。
-
範囲クエリをサポートします。
例えば、クエリ条件が
c < 3の場合、システムはグローバルインデックスに基づいてバケット 2 とバケット 3 を除外し、バケット 0 とバケット 1 からデータをクエリできます。 ハッシュクラスタリングの場合、バケットプルーニングは等価クエリに対してのみ実行できます。 -
複数キー クエリをサポートします。
たとえば、テーブル作成時に
RANGE CLUSTERED BY (c1, c2, c3) SORTED BY (c1, c2, c3)を指定すると、レンジクラスタリングとデータ格納が c1、c2、c3 の順で行われるため、c1 = 100 AND c2 > 0やc1 = 100 AND c2 = 50 AND c3 < 5のような複雑な条件に基づいてテーブルデータにクエリを実行できます。このタイプのクエリは、ハッシュクラスタリングでは実装できません。重要複数キーに基づくクエリでは、クエリ条件のキーを順序どおりに並べる必要があり、値の範囲を指定できるのは最後のキーのみです。
-
レンジクラスタリングは、グローバルソートを効率的に実現します。
レンジクラスタリングを使用する前は、MaxCompute は 1 つのインスタンスのみを使用してデータをグローバルソートすることしかできなかったため、効率が低くなっていました。レンジクラスタリングを使用すると、各範囲内のデータを並列にソートしてから結合できるため、効率が大幅に向上します。
注意事項
レンジクラスタリングの構文はハッシュクラスタリングの構文と似ています。主な違いは、RANGE キーワードを使用する点と、バケット数の指定が任意である点です。
レンジクラスタリングテーブルの作成
CREATE TABLE ステートメントを使用して、レンジクラスタリングテーブルを作成できます。このステートメントでは RANGE CLUSTERED BY パラメーターの指定が必須です。INTO number_of_buckets BUCKETS と SORTED BY パラメーターは任意です。最適な最適化効果を得るため、多くの場合 SORTED BY と RANGE CLUSTERED BY には同じキーを指定することを推奨します。
-
構文
CREATE TABLE [IF NOT EXISTS][(<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] -
例
-
非パーティション化テーブル
CREATE TABLE T1 (a STRING, b STRING, c INT) RANGE CLUSTERED BY (c) SORTED BY (c) INTO 1024 BUCKETS; -
パーティションテーブル
CREATE TABLE T1 (a STRING, b STRING, c INT) PARTITIONED BY (dt INT) RANGE CLUSTERED BY (c) SORTED BY (c) INTO 1024 BUCKETS;
-
-
パラメーター
-
RANGE CLUSTERED BY
レンジクラスタリングのキーを指定します。このパラメーターを指定すると、MaxCompute は指定したバケット数に基づいて、1 つ以上の列のデータを適切な範囲にソートしてサンプリングします。データスキューの問題やホットスポットを防止し、並列クエリのパフォーマンスを向上させるため、
RANGE CLUSTERED BYには、値の範囲が広く、重複キー値が少なく、かつ頻繁に使用される集約キーまたはフィルターキーを指定することを推奨します。 -
SORTED BY
バケット内のフィールドのソート方法を指定します。最適なクエリパフォーマンスを得るため、
SORTED BYとRANGE CLUSTERED BYに同じキーを指定することを推奨します。SORTED BYを指定すると、MaxCompute はグローバルインデックスとファイルインデックスを自動的に生成し、これらのインデックスを使用してクエリを高速化します。 -
INTO number_of_buckets BUCKETS
ハッシュクラスタリングとは異なり、レンジクラスタリングでは
INTO number_of_buckets BUCKETSは省略可能です。 このパラメーターを指定しない場合、MaxCompute はデータ量に基づいてバケット数を自動的に決定します。 ほとんどの場合、実際の状況に基づいてバケット数を指定することをお勧めします。ハッシュクラスタリングと同様に、レンジクラスタリングではバケットサイズ (512 MB ~ 1 GB) に基づいてバケット数を設定することを推奨します。非常に大きなテーブルの場合、多数のバケットが必要です。ただし、テーブル内のバケット数は 4,000 を超えないことを推奨します。
-
テーブルのレンジクラスタリング属性の変更
パーティションテーブルの場合、ALTER 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]]; -- レンジクラスタリングテーブルを非レンジクラスタリングテーブルに変更します。 ALTER TABLE <table_name> NOT CLUSTERED; -
使用上の注意
-
ALTER TABLEステートメントで変更できるのは、パーティションテーブルのクラスタリング属性のみです。非パーティション化テーブルの場合、テーブルにクラスタリング属性を追加した後にクラスタリング属性を変更することはできません。 -
ALTER TABLEステートメントは、テーブルの新しいパーティション (INSERT OVERWRITEステートメントで生成された新しいパーティションを含む) にのみ有効になります。新しいパーティションはクラスタリング属性に基づいて格納されます。既存パーティションのストレージ形式は変更されません。 -
ALTER TABLEステートメントはテーブルの新しいパーティションにのみ有効です。そのため、このステートメントでパーティションを指定することはできません。
-
ALTER TABLE ステートメントは既存テーブルに適しています。レンジクラスタリング属性を追加すると、新しいパーティションはレンジクラスタリング属性に基づいて格納されます。
テーブル属性の明示的な検証
レンジクラスタ化テーブルを作成した後、次のステートメントを実行してテーブルのプロパティを表示できます。 レンジクラスタリングプロパティは、返された結果の Extended Info に表示されます。
DESC EXTENDED <table_name>;
出力例:
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] |
+----------------------------+
Extended Info セクションでは、ClusterType が range、BucketNum が 1024、ClusterColumns が [l_orderkey]、SortColumns が [l_orderkey ASC] となっており、テーブルがレンジクラスタリング属性で正常に構成されていることを確認できます。
パーティションテーブルの場合、次のステートメントを実行してパーティションのクラスタリング属性も確認できます。
DESC EXTENDED <table_name> partition(<pt_spec>);
出力では、ClusterType、BucketNum、ClusterColumns、SortColumns の 4 つの属性が、パーティションのレンジクラスタリング構成を示します。
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]
シナリオ
フィルタリングによるクエリの最適化
テーブルでレンジクラスタリングが有効になっている場合、テーブル内のデータはグローバルにソートされます。MaxCompute は、ソート済みデータに基づいてグローバルインデックスとファイルインデックスを自動的に作成します。これにより、データの格納特性に基づくデータフィルタリングの効率が向上します。レンジクラスタリングを使用して、等価クエリと範囲クエリを最適化できます。
例えば、簡単なクエリ条件 id < 3 の場合、システムはオプティマイザから条件を抽出し、その条件を値の範囲 (-∞, 3) に変換します。 この場合、システムはグローバルインデックスを使用してバケットプルーニングを実行し、前述の値の範囲にデータが含まれないバケット 2 とバケット 3 の両方を除外できます。 次に、システムはバケット 0 とバケット 1 の各ファイルのインデックスを使用して、データを迅速に特定できます。 このプロセスは、次の図に示すように、述語プッシュダウンと呼ばれます。
次のサンプルステートメントは、範囲クラスタリングが実行された 100 GB のデータセットからデータをクエリするための TPC-H クエリ 6 ステートメントです。 TPC-H クエリ 6 ステートメントでは、範囲フィルタリングに基づいて集計操作が実行されます。 範囲クラスタリングは、2 レベルのインデックスを使用してデータを迅速に特定できます。 これにより、クエリの実行時間と消費される CPU およびメモリリソースが大幅に削減されます。
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;
複数キー クエリ
この例では、次のステートメントを使用して mf_tab テーブルを範囲クラスタ化テーブルに変更します。これにより、複数キーのクエリをよりよく理解できます。
ALTER TABLE mf_project.mf_tab
RANGE CLUSTERED BY (project_name, name)
SORTED BY (project_name, name)
INTO 1024 BUCKETS;
テーブルをレンジクラスタリングテーブルに変更した後、プロジェクトレベルで集約クエリを実行できます。サンプルステートメント:
SELECT COUNT(*)
FROM mf_project.mf_tab
WHERE project_name="xxxdw"
AND ds="20180115"
AND type="TABLE";
また、複数のキーを使用してテーブルを正確に特定することもできます。サンプルステートメント:
SELECT COUNT(*)
FROM mf_project.mf_tab
WHERE project_name="xxxdw"
AND name="adm_ctu_cle_kba_midun_trade_dd"
AND type="TABLE";
範囲クエリには複数のキーを使用することもできます。次のステートメントは、名前が adm で始まるテーブルを照会します。
SELECT COUNT(*)
FROM mf_project.mf_tab
WHERE project_name="xxxdw"
AND name >= "adm"
AND name < "adn"
AND type="TABLE";
前述のすべてのクエリは、レンジクラスタリングのグローバルソート機能を最大限に活用し、述語プッシュダウンを実行してテーブルスキャンの I/O 操作数を削減できます。これにより、データのフィルタリングと計算に消費される CPU およびメモリリソースを節約できます。
レンジクラスタリングで複数のキーを使用する場合、特定の要件を満たす必要があります。 テーブル作成ステートメントの RANGE CLUSTERED BY k0, k1, ..., kn において、 km をデータクエリで使用する場合、 k0, k1, ..., km-1 のすべてを条件で指定し、かつそのすべての条件を等式条件にする必要があります。 これにより、インデックスベースのクエリ高速化の最適なパフォーマンスを達成できます。
例えば、k1, k2 は、T という名前のテーブルのクラスターキーです。
-
クエリ条件が
k1 < 5の場合、インデックスに基づくクエリの高速化を実現できます。 -
クエリ条件が
k1 = 10 AND k2 = 20の場合、インデックスによるクエリの高速化が可能です。 -
クエリ条件が
k1 = 10 AND k2 < 0の場合、インデックスによるクエリ高速化を実現できます。 -
クエリ条件が
k2 < 0の場合、インデックスによるクエリの高速化は達成できません。これは、クエリ条件で k1 が指定されていないためです。 -
クエリ条件が
k1 < 0 AND k2 > 0の場合、インデックスを使用したクエリの高速化により、k1 < 0の条件を満たすデータを取得できます。k2 > 0の条件を満たすデータについては、テーブルスキャンが必要です。
GROUP BY の最適化
テーブルでレンジクラスタリングが有効になっている場合、テーブル内のデータはグローバルにソートされます。レンジクラスタリング中、同じ値を持つキーは同じバケットに配置されます。このデータの物理特性を利用することで、集約操作中のシャッフル手順を省略できます。
たとえば、次の CREATE TABLE ステートメントを使用して T という名前のテーブルを作成します。このテーブルからデータをクエリする場合、map ステージでテーブルデータに対して GROUP BY 操作を実行できます。
CREATE TABLE T (department INT, team STRING, employee STRING)
RANGE CLUSTERED BY (department, team)
SORTED BY (department, team)
INTO 1024 BUCKETS;
SELECT COUNT(*) FROM T GROUP BY department, team;
GROUP BY の最適なパフォーマンスを得るには、GROUP BY と RANGE CLUSTERED BY に同じキーを指定する必要があります。
集約の最適化
次のステートメントは、テーブル foo のデータ構造を示します。
CREATE TABLE foo(a BIGINT, b BIGINT, c BIGINT)
RANGE CLUSTERED BY (a,b)
SORTED BY(a,b) INTO 3 BUCKETS;
テーブル foo のバケットに格納されているデータは、以下の範囲に分類されます。
Bucket 0: [1,1 : 3,3]
Bucket 1: [5,5 : 7,7]
Bucket 2: [8,8 : 9,9]
前述のバケット範囲は、Bucket N: [lower bound values : upper bound values] の形式で指定されます。データが列 a で集計される場合、代わりに次のバケット範囲が使用されます。
Bucket 0: [1 : 3]
Bucket 1: [5 : 7]
Bucket 2: [8 : 9]
列 a と列 b による集約操作の実行計画を直接生成し、3 つのインスタンスを起動して各バケット内のデータを集約し、出力結果を返すことができます。
ただし、列 a の値が複数のバケットに分散している場合は、正しくない結果が返されます。例:
Bucket 0: [1,1 : 3,3]
Bucket 1: [3,5 : 7,7]
Bucket 2: [7,8 : 9,9]
列 a には値 3 と 7 があり、これらは 2 つのバケットに別々に格納されます。有効な結果を得るには、列 a で同じ値を持つタプルを同じインスタンスに配置し、タプルを集計する必要があります。これにより、次の図に示すように、バケットが再度作成されます。赤色の 2 本の破線間のスペースは、各インスタンスが読み取り可能なデータ範囲を指定します。
範囲クラスタリングにはヒストグラムが必要です。範囲クラスター化テーブルでは、クラスターキーとソートキーが同じ場合、テーブルにデータを挿入する際に、各バケットに対応するワーカーが 10,000 行ごとにタプルをサンプリングし、クラスターキーの値を取得します。取得された値はヒストグラムに保存されます。各バケットのヒストグラムは、クラスターメタデータファイルに格納されます。このタイプのヒストグラムは、等深度ヒストグラムと呼ばれます。
タプルのサンプリングは、クラスターキーとソートキーが同一の場合にのみ行われます。
各バケットのヒストグラムを取得した後、次のルールに基づいて各ワーカーのバケットを再作成できます。
-
同じグルーピングキーを持つタプルは同じバケットに格納されます。
-
データはバケット間で均等に分布します。
各新規バケットの下限値に基づき、各ワーカーは有効な範囲内のデータを読み取り、正しい結果を返すことができます。
以下では、TPC-H データセット内の 1 TB のデータを持つ partsupp テーブルを使用して、パフォーマンスの向上をテストします。次のステートメントを実行して、partsupp テーブルをレンジクラスタリングテーブルに変換します。
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;
次のクエリステートメントを実行してテストを行います。
SELECT ps_partkey, count(*) c FROM partsupp GROUP BY ps_partkey;
-
次のコマンドを実行して最適化を無効にします。
SET odps.optimizer.enable.range.partial.repartitioning=false;以下に、最適化が無効な場合のジョブ実行ログを示します。
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) -
次のコマンドを実行して最適化を有効にします。
SET odps.optimizer.enable.range.partial.repartitioning=true;出力例:
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)
テスト結果から、最適化を有効にすると、クエリ速度が 57% 向上し、CPU 使用率が 52% 低下し、メモリ使用量が 71% 低下することがわかります。パフォーマンスの向上は、データ量とクエリタイプによって異なります。
レンジクラスタリングテーブルの結合最適化
-
この例では、次のステートメントを使用して 2 つのテーブルを作成します。
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;その後、2 つのテーブルに異なるデータを挿入します。
結合が必要な 2 つのハッシュクラスタリングテーブルの場合、両テーブルのバケット数が同じであれば、テーブルのバケット内のデータを結合できます。ただし、このルールはレンジクラスタリングテーブルには適用されません。結合が必要な 2 つのレンジクラスタリングテーブルでは、両テーブルのバケット数が同じであっても、バケット ID に基づいてバケット内のデータを直接結合できません。これは、レンジクラスタリングテーブルの各バケット境界が異なる場合があるためです。2 つのレンジクラスタリングテーブルを結合すると、次の図に示すように、シャッフル手順を含む実行計画が常に生成されます。
2 つのレンジクラスタリングテーブル間の結合を最適化するには、テーブルの境界を揃えて 2 つのテーブルのバケットを再作成します。これにより、各インスタンスが読み取れるデータ境界が再定義されます。 -
2 つのテーブルを作成します。
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;テーブルに一定量のデータを挿入すると、次の図に示すようにバケット境界が定義されます。
サンプルクエリ 1:SELECT * FROM t1 JOIN t2 ON t1.a=t2.a AND t1.b=t2.b AND t1.c=t2.c;オプティマイザーは、バケット数が多いテーブルの境界をもう一方のテーブルの境界に揃え、次の図に示すように各テーブルの新しい境界を取得します。
これにより、シャッフル手順を含まない実行計画が生成されます。
サンプルクエリ 2:SELECT * FROM t1 JOIN t2 ON t1.a=t2.a AND t1.b=t2.b;オプティマイザーは列 a と列 b に基づいて各テーブルのバケットを作成し、その後、整列してバケットを再作成することで、各テーブルの各バケットからデータを読み取る境界を取得します。これにより、前述の図に示すように、シャッフル手順を含まない実行計画が生成されます。
-
パフォーマンステスト
-
テーブル変換
TPC-H Query 2 ステートメントを使用して、PART および PARTSUPP という名前の 2 つのテーブルでテストを行います。各テーブルには 1 TB のデータが含まれます。2 つのテーブルをレンジクラスタリングテーブルに変換し、他のテーブルは変更しません。
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;次の TPC-H Query 2 ステートメントを使用してデータをクエリします。
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; -
テスト結果
-
次のコマンドを実行して最適化を無効にします。
SET odps.optimizer.enable.range.partial.repartitioning=false;以下に、最適化が無効な場合のジョブ実行ログ出力例を示します。
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) -
次のコマンドを実行して最適化を有効にします。
SET odps.optimizer.enable.range.partial.repartitioning=true;出力例:
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)
最適化後、2 つのステージが削除され、クエリ速度は約 21.4% 向上し、CPU 使用率は約 35.4% 低下し、メモリ使用量は約 54.6% 低下します。
-
-
グローバルソートの高速化
レンジクラスタリングは、グローバルソートの高速化にも使用できます。一般的な ORDER BY のシナリオでは、グローバルソートを保証するために、ソート済みデータが同じインスタンスに分散されます。しかし、これらのシナリオでは並列処理を十分に活用できません。レンジクラスタリングのパーティショニング手順を使用して、並列グローバルソートを実装できます。グローバルソートでは、データをサンプリングして範囲に分割し、各範囲内のデータを並列にソートした後、グローバルソートの結果を取得する必要があります。
グローバルソートが完了した後も、テーブルまたはテーブル内のパーティションのクラスター属性を変更すると、テーブルには複数のバケットが引き続き含まれます。データを利用する際には、グローバルソートを保証するために、バケット ID に基づいてファイル内のデータを読み取る必要があります。
デフォルトでは、レンジクラスタリングテーブルのグローバルソート高速化は無効です。グローバルソート高速化を有効にするには、次のコマンドを実行します。
SET odps.optimizer.distribute.ordering.enable=true;
制限事項と使用上の注意
ハッシュクラスタリングと比較して、レンジクラスタリングには次の制限があります。
-
レンジクラスタリングのデータ生成コストは、ハッシュクラスタリングより高くなります。ハッシュクラスタリングは、データのハッシュ化とソートという単純な操作のみを必要とします。一方、レンジクラスタリングでは、データのサンプリング、ソート、ヒストグラムの結合が必要です。実行時間、CPU コスト、メモリコストを含む全体的な消費は、ハッシュクラスタリングよりも高くなります。そのため、ハッシュクラスタリングで問題を解決できる場合は、レンジクラスタリングを使用する必要はありません。
-
レンジクラスタリングは
DYNAMIC PARTITIONまたはINSERT INTOではサポートされていません。 -
レンジクラスタリングでサポートされる結合操作は、内部結合、左外部結合、右外部結合、セミ結合のみです。レンジクラスタリングはアンチ結合または完全外部結合ではサポートされていません。
-
レンジクラスター化テーブルでは、RANGE CLUSTERED BY で指定するキーは、SORTED BY で指定するキーと同じである必要があります。 たとえば、テーブル作成時に foo という名前のテーブルに
range clustered by (a,b) sorted by (a,b)を指定した場合、このトピックで説明している最適化をそのテーブルで実現できます。 しかし、bar という名前のテーブルにrange clustered by(a,b) sorted by (b,a)を指定した場合、このトピックで説明している最適化をそのテーブルで実現することはできません。 -
JOIN または GROUP BY で指定されるキーは、RANGE CLUSTERED BY で指定されるキーのプレフィックスまたはすべてのキーである必要があります。 たとえば、テーブル作成ステートメントで
range clustered by(a,b,c) sorted by(a,b,c)が指定されているとします。 このトピックで説明されている最適化は、JOIN または GROUP BY のキーとしてa、a,b、またはa,b,cが指定されている場合にのみ、テーブルで実現できます。 JOIN または GROUP BY のキーとしてbまたはa,cが指定されている場合、このトピックで説明されている最適化は実現できません。 -
レンジクラスタリングされたパーティションテーブルでは、テーブルの 2 つ以上のパーティションからデータを読み取る場合、このトピックで説明する最適化を実現できません。このトピックで説明する最適化を実現できるのは、単一パーティションを持つパーティションテーブルと非パーティション化テーブルのみです。