OceanBaseデータベースはHash Join結合アルゴリズムをサポートしており、特定のフィールドに基づいて2つのテーブルを等価マッチで結合することができます。しかし、結合に参加するテーブル、特にProbe Tableのデータ量が大きい場合、Hash Joinのパフォーマンスは大幅に低下します。このような状況では、Hash JoinはRuntime Filter(RF)を使用して効率を向上させることができます。
原理の紹介
Runtime Filterは、Hash Joinのパフォーマンスを最適化するための技術であり、Hash Joinの使用時に必要なProbeのデータ量を削減することでクエリの効率を向上させます。例えば、スター型結合の複数の次元テーブルと実装テーブルをJoinするシナリオでは、Runtime Filterは非常に有効な最適化手法です。
Runtime Filterは実際にはフィルターであり、Hash JoinのBuildプロセスを利用して軽量なFilterを構築し、そのFilterをProbeに参加する大規模テーブルにブロードキャストします。Probe Tableは複数のRuntime Filterを使用してストレージ層で事前にデータをフィルタリングし、実際にHash Joinおよびネットワーク転送に参加するデータ量を削減することで、クエリの効率を向上させます。
Runtime Filterは3つの観点から分類できます。マシントランスポートが必要かどうかに基づいて、Runtime FilterはLocalとGlobalの2種類に分けられます。Runtime Filterのデータ構造に基づいて、Runtime FilterはBloom Filter、In Filter、Range Filterの3種類に分けられます。Local Runtime FilterであってもGlobal Runtime Filterであっても、これら3種類のデータ構造を持つRuntime Filterを使用できます。Runtime Filterがフィルタリングするエンティティに基づいて、Runtime Filterは接続キーのフィルタリングとパーティションのフィルタリングの2種類に分けられます。前者が一般的なRuntime Filterであり、後者がPart Join Filterです。フィルタリングするエンティティに基づく分類とマシントランスポートが必要かどうかに基づく分類は直交しています。一般的なRuntime Filterは3種類のデータ構造を持つRuntime Filterをサポートしますが、Part Join Filterは現在Bloom Filterのみをサポートしています。
本記事では、具体的な例を用いてLocal Runtime Filter、Global Runtime Filter、Part Join Filterについて詳しく説明します。
Local Runtime Filter
Local Runtime FilterのRuntime Filterはネットワーク転送を経由する必要がなく、構築されたFilterはローカルノードでフィルタリング条件を計算するだけで済みます。Local Runtime Filterは通常、Hash JoinのProbe側にShuffleがないシナリオに適用されます。
以下はLocal Runtime Filterの実行計画の例です。
obclient> CREATE TABLE tt1(v1 INT, v2 INT) PARTITION BY HASH(v1) PARTITIONS 5;
Query OK, 0 rows affected
obclient> CREATE TABLE tt2(v1 INT, v2 INT) PARTITION BY HASH(v1) PARTITIONS 5;
Query OK, 0 rows affected
obclient> EXPLAIN SELECT /*+ PX_JOIN_FILTER(tt2) PARALLEL(4) */ * FROM tt1 JOIN tt2 ON tt1.v1=tt2.v1;
+------------------------------------------------------------------------------------------------------------------------------------+
| Query Plan |
+------------------------------------------------------------------------------------------------------------------------------------+
| ============================================================== |
| |ID|OPERATOR |NAME |EST.ROWS|EST.TIME(us)| |
| -------------------------------------------------------------- |
| |0 |PX COORDINATOR | |1 |7 | |
| |1 |└─EXCHANGE OUT DISTR |:EX10000|1 |6 | |
| |2 | └─PX PARTITION ITERATOR | |1 |6 | |
| |3 | └─HASH JOIN | |1 |6 | |
| |4 | ├─JOIN FILTER CREATE|:RF0000 |1 |3 | |
| |5 | │ └─TABLE FULL SCAN |tt2 |1 |3 | |
| |6 | └─JOIN FILTER USE |:RF0000 |1 |4 | |
| |7 | └─TABLE FULL SCAN |tt1 |1 |4 | |
| ============================================================== |
| Outputs & filters: |
| ------------------------------------- |
| 0 - output([INTERNAL_FUNCTION(tt1.v1, tt1.v2, tt2.v1, tt2.v2)]), filter(nil), rowset=256 |
| 1 - output([INTERNAL_FUNCTION(tt1.v1, tt1.v2, tt2.v1, tt2.v2)]), filter(nil), rowset=256 |
| dop=4 |
| 2 - output([tt1.v1], [tt2.v1], [tt2.v2], [tt1.v2]), filter(nil), rowset=256 |
| partition wise, force partition granule |
| 3 - output([tt1.v1], [tt2.v1], [tt2.v2], [tt1.v2]), filter(nil), rowset=256 |
| equal_conds ([tt1.v1 = tt2.v1]), other_conds(nil) |
| 4 - output([tt2.v1], [tt2.v2]), filter(nil), rowset=256 |
| RF_TYPE(in, range, bloom), RF_EXPR[tt2.v1] |
| 5 - output([tt2.v1], [tt2.v2]), filter(nil), rowset=256 |
| access([tt2.v1], [tt2.v2]), partitions(p[0-4]) |
| is_index_back=false, is_global_index=false, |
| range_key([tt2.__pk_increment]), range(MIN ; MAX)always true |
| 6 - output([tt1.v1], [tt1.v2]), filter(nil), rowset=256 |
| 7 - output([tt1.v1], [tt1.v2]), filter([RF_IN_FILTER(tt1.v1)], [RF_RANGE_FILTER(tt1.v1)], [RF_BLOOM_FILTER(tt1.v1)]), rowset=256 |
| access([tt1.v1], [tt1.v2]), partitions(p[0-4]) |
| is_index_back=false, is_global_index=false, filter_before_indexback[false,false,false], |
| range_key([tt1.__pk_increment]), range(MIN ; MAX)always true |
+------------------------------------------------------------------------------------------------------------------------------------+
32 rows in set
上記の例では、NAME フィールドの値が RF0000 である4番目の JOIN FILTER CREATE と6番目の JOIN FILTER USE 演算子が一般的なRuntime Filterです。これらは計画内でLocal Runtime Filterを構築し、このFilterはネットワーク伝播を必要とせず、ローカルでのみ使用されます。
Global Runtime Filter
Global Runtime FilterのRuntime Filterは複数の実行ノードにブロードキャストされ、計画の形態では必要に応じてRFをネストして計画内の任意の位置まで押し下げてフィルタリングを完了できます。Local Runtime Filterと比較して、Global Runtime Filterによるコスト削減は、SQL層の投影および計算オーバーヘッドだけでなく、ネットワーク転送も含まれるため、通常は良好な実行性能向上が得られます。
以下はGlobal Runtime Filterの実行計画の例です。
obclient> CREATE TABLE tt1 (c1_rand INT, c2_rand INT, c3_rand INT, c4_rand INT, c5_rand INT) PARTITION BY HASH(c5_rand) PARTITIONS 5;
Query OK, 0 rows affected
obclient> CREATE TABLE tt2 (c1_rand INT, c2_rand INT, c3_rand INT, c4_rand INT, c5_rand INT) PARTITION BY HASH(c5_rand) PARTITIONS 5;
Query OK, 0 rows affected
obclient> CREATE TABLE tt3 (c1_rand INT, c2_rand INT, c3_rand INT, c4_rand INT, c5_rand INT) PARTITION BY HASH(c5_rand) PARTITIONS 5;
Query OK, 0 rows affected
obclient> EXPLAIN BASIC SELECT /*+ LEADING(a (b c)) PARALLEL(3) USE_HASH(b) USE_HASH(c) PQ_DISTRIBUTE(c HASH HASH) PX_JOIN_FILTER(c a) PX_JOIN_FILTER(c b) */ COUNT(*) FROM tt1 a, tt2 b, tt3 c WHERE a.c1_rand=b.c1_rand AND a.c2_rand = c.c2_rand AND b.c3_rand = c.c3_rand;
+------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| Query Plan |
+------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| ======================================================== |
| |ID|OPERATOR |NAME | |
| -------------------------------------------------------- |
| |0 |SCALAR GROUP BY | | |
| |1 |└─PX COORDINATOR | | |
| |2 | └─EXCHANGE OUT DISTR |:EX10003| |
| |3 | └─MERGE GROUP BY | | |
| |4 | └─SHARED HASH JOIN | | |
| |5 | ├─JOIN FILTER CREATE |:RF0001 | |
| |6 | │ └─EXCHANGE IN DISTR | | |
| |7 | │ └─EXCHANGE OUT DISTR (BC2HOST)|:EX10000| |
| |8 | │ └─PX BLOCK ITERATOR | | |
| |9 | │ └─TABLE FULL SCAN |a | |
| |10| └─HASH JOIN | | |
| |11| ├─JOIN FILTER CREATE |:RF0000 | |
| |12| │ └─EXCHANGE IN DISTR | | |
| |13| │ └─EXCHANGE OUT DISTR (HASH) |:EX10001| |
| |14| │ └─PX BLOCK ITERATOR | | |
| |15| │ └─TABLE FULL SCAN |b | |
| |16| └─EXCHANGE IN DISTR | | |
| |17| └─EXCHANGE OUT DISTR (HASH) |:EX10002| |
| |18| └─JOIN FILTER USE |:RF0000 | |
| |19| └─JOIN FILTER USE |:RF0001 | |
| |20| └─PX BLOCK ITERATOR | | |
| |21| └─TABLE FULL SCAN |c | |
| ======================================================== |
| Outputs & filters: |
| ------------------------------------- |
| 0 - output([T_FUN_COUNT_SUM(T_FUN_COUNT(*))]), filter(nil), rowset=256 |
| group(nil), agg_func([T_FUN_COUNT_SUM(T_FUN_COUNT(*))]) |
| 1 - output([T_FUN_COUNT(*)]), filter(nil), rowset=256 |
| 2 - output([T_FUN_COUNT(*)]), filter(nil), rowset=256 |
| dop=3 |
| 3 - output([T_FUN_COUNT(*)]), filter(nil), rowset=256 |
| group(nil), agg_func([T_FUN_COUNT(*)]) |
| 4 - output(nil), filter(nil), rowset=256 |
| equal_conds([a.c1_rand = b.c1_rand], [a.c2_rand = c.c2_rand]), other_conds(nil) |
| 5 - output([a.c2_rand], [a.c1_rand]), filter(nil), rowset=256 |
| RF_TYPE(in, range, bloom), RF_EXPR[a.c2_rand] |
| 6 - output([a.c2_rand], [a.c1_rand]), filter(nil), rowset=256 |
| 7 - output([a.c2_rand], [a.c1_rand]), filter(nil), rowset=256 |
| dop=3 |
| 8 - output([a.c1_rand], [a.c2_rand]), filter(nil), rowset=256 |
| 9 - output([a.c1_rand], [a.c2_rand]), filter(nil), rowset=256 |
| access([a.c1_rand], [a.c2_rand]), partitions(p[0-4]) |
| is_index_back=false, is_global_index=false, |
| range_key([a.__pk_increment]), range(MIN ; MAX)always true |
| 10 - output([b.c1_rand], [c.c2_rand]), filter(nil), rowset=256 |
| equal_conds([b.c3_rand = c.c3_rand]), other_conds(nil) |
| 11 - output([b.c3_rand], [b.c1_rand]), filter(nil), rowset=256 |
| RF_TYPE(in, range, bloom), RF_EXPR[b.c3_rand] |
| 12 - output([b.c3_rand], [b.c1_rand]), filter(nil), rowset=256 |
| 13 - output([b.c3_rand], [b.c1_rand]), filter(nil), rowset=256 |
| (#keys=1, [b.c3_rand]), dop=3 |
| 14 - output([b.c1_rand], [b.c3_rand]), filter(nil), rowset=256 |
| 15 - output([b.c1_rand], [b.c3_rand]), filter(nil), rowset=256 |
| access([b.c1_rand], [b.c3_rand]), partitions(p[0-4]) |
| is_index_back=false, is_global_index=false, |
| range_key([b.__pk_increment]), range(MIN ; MAX)always true |
| 16 - output([c.c3_rand], [c.c2_rand]), filter(nil), rowset=256 |
| 17 - output([c.c3_rand], [c.c2_rand), filter(nil), rowset=256 |
| (#keys=1, [c.c3_rand]), dop=3 |
| 18 - output([c.c3_rand], [c.c2_rand]), filter(nil), rowset=256 |
| 19 - output([c.c3_rand], [c.c2_rand]), filter(nil), rowset=256 |
| 20 - output([c.c3_rand], [c.c2_rand]), filter(nil), rowset=256 |
| 21 - output([c.c3_rand], [c.c2_rand]), filter([RF_IN_FILTER(c.c3_rand)], [RF_RANGE_FILTER(c.c3_rand)], [RF_BLOOM_FILTER(c.c3_rand)], [RF_IN_FILTER(c.c2_rand)], |
| [RF_RANGE_FILTER(c.c2_rand)], [RF_BLOOM_FILTER(c.c2_rand)]), rowset=256 |
| access([c.c2_rand], [c.c3_rand]), partitions(p[0-4]) |
| is_index_back=false, is_global_index=false, filter_before_indexback[false,false,false,false,false,false], |
| range_key([c.__pk_increment]), range(MIN ; MAX)always true |
+------------------------------------------------------------------------------------------------------------------------------------------------------------------+
70 rows in set
上記の例では、5番目の演算子の NAME フィールドの値が RF0001 に対応し、11番目の演算子が RF000 に対応します。そして、21番目にPush Downされた TABLE FULL SCAN 演算子は複数のDFOにまたがってフィルタリングを行い、Global Runtime Filterの特徴に合致しています。
Part Join Filter
Hash Joinの実行プロセスは、左側のデータでハッシュテーブルを構築し、右側で行ごとにデータをマッチングするものです。左側でハッシュテーブルを構築する際には、左側のすべてのデータを取得します。このプロセスで、もし左側のすべてのデータに関する右側の特定のテーブルのパーティション分布特性を取得できれば、右側でそのテーブルのデータをスキャンする際に、既に統計されたパーティション分布特性に基づいて不要なパーティションを事前にフィルタリングすることができ、パフォーマンスを向上させることができます。Part Join Filterの導入は、まさにこのシナリオを最適化するためのものです。Hash Joinの左側で右側の特定のテーブルの具体的なパーティションを計算するためには、Joinの接続キーに右側のそのテーブルのパーティションキーを含める必要があります。これがPart Join Filterを生成する前提条件です。
以下はPart Join Filterの実行計画の例です。
obclient> CREATE TABLE tt1(v1 INT);
Query OK, 0 rows affected
obclient> CREATE TABLE tt2(v1 INT) PARTITION BY HASH(v1) PARTITIONS 5;
Query OK, 0 rows affected
obclient> EXPLAIN SELECT /*+ PARALLEL(3) LEADING(tt1 tt2) PX_PART_JOIN_FILTER(tt2)*/ * FROM tt1 JOIN tt2 ON tt1.v1=tt2.v1;
+-------------------------------------------------------------------------------+
| Query Plan |
+-------------------------------------------------------------------------------+
| ======================================================================= |
| |ID|OPERATOR |NAME |EST.ROWS|EST.TIME(us)| |
| ----------------------------------------------------------------------- |
| |0 |PX COORDINATOR | |1 |5 | |
| |1 |└─EXCHANGE OUT DISTR |:EX10001|1 |4 | |
| |2 | └─HASH JOIN | |1 |4 | |
| |3 | ├─PART JOIN FILTER CREATE |:RF0000 |1 |1 | |
| |4 | │ └─EXCHANGE IN DISTR | |1 |1 | |
| |5 | │ └─EXCHANGE OUT DISTR (PKEY)|:EX10000|1 |1 | |
| |6 | │ └─PX BLOCK ITERATOR | |1 |1 | |
| |7 | │ └─TABLE FULL SCAN |tt1 |1 |1 | |
| |8 | └─PX PARTITION HASH JOIN-FILTER|:RF0000 |1 |3 | |
| |9 | └─TABLE FULL SCAN |tt2 |1 |3 | |
| ======================================================================= |
| Outputs & filters: |
| ------------------------------------- |
| 0 - output([INTERNAL_FUNCTION(tt1.v1, tt2.v1)]), filter(nil), rowset=256 |
| 1 - output([INTERNAL_FUNCTION(tt1.v1, tt2.v1)]), filter(nil), rowset=256 |
| dop=3 |
| 2 - output([tt1.v1], [tt2.v1]), filter(nil), rowset=256 |
| equal_conds([tt1.v1 = tt2.v1]), other_conds(nil) |
| 3 - output([tt1.v1]), filter(nil), rowset=256 |
| RF_TYPE(bloom), RF_EXPR[t1.v1] |
| 4 - output([tt1.v1]), filter(nil), rowset=256 |
| 5 - output([tt1.v1]), filter(nil), rowset=256 |
| (#keys=1, [t1.v1]), dop=3 |
| 6 - output([tt1.v1]), filter(nil), rowset=256 |
| 7 - output([tt1.v1]), filter(nil), rowset=256 |
| access([tt1.v1]), partitions(p0) |
| is_index_back=false, is_global_index=false, |
| range_key([tt1.__pk_increment]), range(MIN ; MAX)always true |
| 8 - output([tt2.v1]), filter(nil), rowset=256 |
| affinitize |
| 9 - output([tt2.v1]), filter(nil), rowset=256 |
| access([tt2.v1]), partitions(p[0-4]) |
| is_index_back=false, is_global_index=false, |
| range_key([tt2.__pk_increment]), range(MIN ; MAX)always true |
+-------------------------------------------------------------------------------+
上記の例では、計画内の3番目の PART JOIN FILTER CREATE 演算子と8番目の PX PARTITION HASH JOIN-FILTER 演算子がPart Join Filterです。3番目の演算子がPart Join Filterに対応し、8番目の演算子で tt2 テーブルに対するパーティションレベルのフィルタリングに適用されます。
Runtime Filterの手動での有効化と無効化
Runtime Filterの使用シナリオはHash Joinに限定されます。接続タイプがHash Join以外の場合、オプティマイザーはRuntime Filterを割り当てません。通常、Hash Joinを実行する際にはオプティマイザーが自動的にRuntime Filterを割り当てますが、オプティマイザーがRuntime Filterを割り当てていない場合、ユーザーはHintを使用して手動で割り当てることもできます。
PX_JOIN_FILTER HintとPX_PART_JOIN_FILTER Hintは、Runtime Filterを手動で有効にするために使用されます。SQL構文は以下のとおりです:
/*+ PX_JOIN_FILTER(join_right_table_name)*/
/*+ PX_PART_JOIN_FILTER(join_right_table_name)*/
例:
EXPLAIN SELECT /*+ PX_JOIN_FILTER(tt2) PARALLEL(4) */ * FROM tt1 JOIN tt2 ON tt1.v1=tt2.v1;
EXPLAIN SELECT /*+ PARALLEL(3) LEADING(tt1 tt2) PX_PART_JOIN_FILTER(tt2)*/ * FROM tt1 JOIN tt2 ON tt1.v1=tt2.v1;
並列度が1のシナリオでは、ランタイムフィルターは割り当てられません。この場合、ユーザーは並列度を少なくとも2に指定する必要があります。PARALLEL ヒントを使用して設定できます。例:
/*+ PARALLEL(2) */
注意点として、Hash Joinのフィルタリング性能が不十分な場合、ランタイムフィルターを使用しても問題は解決せず、むしろわずかなパフォーマンス低下を招く可能性があります。そのため、手動でランタイムフィルターを有効にする際には、クエリシナリオを慎重に評価し、ランタイムフィルターが適用可能かどうかを判断する必要があります。
NO_PX_JOIN_FILTER ヒントと NO_PX_PART_JOIN_FILTER ヒントは、ランタイムフィルターを手動で無効にするために使用されます。SQL構文は以下のとおりです:
/*+ NO_PX_JOIN_FILTER(join_right_table_name)*/
/*+ NO_PX_PART_JOIN_FILTER(join_right_table_name)*/
ランタイムフィルターの実行戦略の調整
ランタイムフィルターは自己適応機能を備えています。デフォルト設定では、Joinが接続キーの選択度の条件を満たすと、In、Range、Bloomの3種類のデータ構造を持つランタイムフィルターがデフォルトで作成されます。In Filterは内部でHashテーブルを使用してフィルタリング判断を行い、Range Filterは内部で最大または最小値を使用してフィルタリング判断を行います。
フィルタリングの優先順位については、In Filterが最も正確なフィルタリング性能を持っているため、実行時にIn Filterが有効になると、他の2種類のFilterは自己適応的に無効になります。実行時には、エグゼキューターが実際のNDV値に基づいてIn Filterを使用するかどうかを判断します。さらに、各Filterは実際の計算において、リアルタイムのフィルタリング性能に応じて自己適応的にDisable FilterおよびReEnable Filterを行います。
OceanBaseデータベースは、runtime_filter_type、runtime_filter_max_in_num、runtime_filter_wait_time_ms、runtime_bloom_filter_max_size を含む、ランタイムフィルターの実行に関する戦略を調整するためのシステム変数を提供しています。
runtime_filter_type 変数は、有効にするランタイムフィルターのタイプを指定するために使用されます。runtime_filter_type のデフォルト値は 'BLOOM_FILTER,RANGE,IN' で、3種類のランタイムフィルターを同時に有効にすることを意味します。runtime_filter_type='' の場合、どのタイプのランタイムフィルターも有効にしません。例:
ALTER SYSTEM SET runtime_filter_type = 'BLOOM_FILTER,RANGE,IN'
通常、デフォルト値で指定されたランタイムフィルターのタイプを使用するだけで済み、runtime_filter_type 変数で特別に指定する必要はありません。OceanBaseデータベースは、最適化および実行段階で最適なランタイムフィルターのタイプを選択してフィルタリングを行います。
runtime_filter_max_in_num 変数は、In Filterが使用するNDVを指定するために使用されます。デフォルト値は1024です。オプティマイザーはBuildテーブルのNDVを推定します。NDV > runtime_filter_max_in_num の場合、In Filterは割り当てられません。オプティマイザーによるNDVの推定は必ずしも正確ではないため、実際のNDVが高いにもかかわらず推定値が低く、誤ってIn Filterが割り当てられる場合があります。このような場合、エグゼキューターは実行中に実際のNDVに基づいてIn Filterを自己適応的に無効にします。通常、この値を大きく設定することは推奨されません。BuildテーブルのNDVが非常に高い場合、In Filterの最適化効果はBloom Filterほど良くありません。
runtime_filter_wait_time_ms 変数は、Probe側がランタイムフィルターの到着を待機する最大時間を設定するために使用されます。デフォルト値は10msです。Probe側はランタイムフィルターが到着するのを待ってからデータを吐き出します。runtime_filter_wait_time_ms で設定された時間までにランタイムフィルターが到着しない場合、By Pass段階に入り、ランタイムフィルターを経由せずに直接データの吐き出しが始まります。後のある時点でランタイムフィルターが到着すると、ランタイムフィルターは有効になり、フィルタリングが実行されます。通常、この値は調整不要です。実際に使用されるBloom Filterのデータが非常に大きく、かつフィルタリング性能が良好な場合には、この値を適切に大きくすることができます。
runtime_bloom_filter_max_size 変数は、Bloom Filterの最大長を制限するために使用されます。デフォルト値は2048MBです。ユーザーが実際に使用する中でBuildテーブルのデータが多すぎる場合、Bloom Filterのデフォルトの最大長ではデータを収容できない場合があります。この場合、runtime_bloom_filter_max_size の値を大きくする必要があります。