複数テーブルの結合クエリでは、結合条件に一致しないデータが I/O オーバーヘッドを増加させ、クエリパフォーマンスを低下させます。複数テーブルの結合シナリオにおいて、ランタイムフィルターはデータスキャンフェーズで一致しないデータをプルーニング (刈り込み) する軽量なフィルターを自動的に生成します。これにより、I/O オーバーヘッドが削減され、クエリパフォーマンスが向上します。この機能は Hologres V2.0 から Hash Join をサポートしており、V4.2 では Cross Join のシナリオにも拡張されました。
背景情報
利用シーン
Hologres は V2.0 からランタイムフィルターをサポートしています。この機能は、2 つ以上のテーブルが関わる Hash Join のシナリオで一般的に使用され、特に大規模テーブルと小規模テーブルを結合する場合に効果的です。手動での構成は不要です。オプティマイザーと実行エンジンがクエリ時に結合フィルタリングの動作を自動的に最適化し、I/O オーバーヘッドを削減して結合パフォーマンスを向上させます。
V4.2 から、ランタイムフィルターの機能は Cross Join のシナリオにも拡張されました。具体的には、スカラーサブクエリで集計結果を計算し、それを大規模テーブルのフィルター条件として使用する一般的な SQL パターンを最適化します。詳細については、「Cross Join のランタイムフィルターサポート (V4.2 の新機能)」をご参照ください。
仕組み
Hash Join のランタイムフィルター
2 つのテーブルが結合される際、一方のテーブルのデータがハッシュテーブルにロードされ、もう一方のテーブルのデータがそれと照合されます。結合プロセスには 2 つの側面があります:
-
ビルド側:ハッシュテーブルを構築する側で、実行計画の Hash ノードに対応します。
-
プローブ側:データを読み取り、ビルド側のハッシュテーブルと照合する側です。
通常、小規模テーブルがビルド側、大規模テーブルがプローブ側として機能します。
ランタイムフィルターは、ビルド側のデータ分布から軽量なフィルターを構築し、それをプローブ側にプッシュしてデータをプルーニングすることで機能します。これにより、Hash Join で処理されるプローブ側のデータ量が削減され、ネットワークトラフィックが最小限に抑えられるため、結合パフォーマンスが向上します。その結果、この機能は、サイズ差が大きい大規模テーブルと小規模テーブル間の結合で最も効果を発揮し、標準的な結合よりも大きなパフォーマンス向上を実現します。
Cross Join のランタイムフィルター (V4.2 の新機能)
V4.2 以前は、ランタイムフィルターは Hash Join のみを対象としていました。しかし、ユーザーはしばしばスカラーサブクエリを使用して min や max などの集計値を計算し、その結果を大規模テーブルのフィルター条件として使用します。実行計画では、この SQL パターンは Cross Join を生成します:
-
ビルド側:スカラーサブクエリの結果で、正確に 1 行です。
-
プローブ側:大規模テーブルのスキャンです。
V4.2 では、新しいフィルタータイプである ScalarFilter が導入されました。オプティマイザーはこのパターンを自動的に識別し、ビルド側で評価されたスカラー値をプローブ側の ScanNode にプッシュダウンします。これにより、スキャンフェーズでの行レベルおよび RowGroup レベルのフィルタリングが可能になり、全表スキャンとその後の行ごとの比較を回避できます。
制限事項とトリガー条件
制限事項
-
ランタイムフィルターは Hologres V2.0 以降でのみサポートされます。
-
Hash Join のシナリオでは、V2.0 は結合条件に単一のフィールドが含まれる場合にのみランタイムフィルターをサポートします。V2.1 以降、ランタイムフィルターは複数のフィールドをサポートします。
-
TopN ランタイムフィルターは V4.0 以降でのみサポートされ、単一テーブルの TopN 計算のパフォーマンスを向上させるために使用されます。
-
Cross Join ランタイムフィルター (ScalarFilter) は V4.2 以降でのみサポートされます。
トリガー条件
Hash Join のシナリオ
エンジンは、以下のすべての条件が満たされた場合にランタイムフィルターを自動的にトリガーします:
-
プローブ側の行数が 100,000 行以上である。
-
ビルド側からスキャンされたデータとプローブ側からスキャンされたデータの比率が 0.1 以下である。この比率が低いほど、フィルターがトリガーされる可能性が高くなります。
-
結合の出力データとプローブ側のデータの比率が 0.1 以下である。この比率が低いほど、フィルターがトリガーされる可能性が高くなります。
Cross Join のシナリオ (V4.2 以降)
オプティマイザーは、以下の条件が満たされた場合に ScalarFilter を自動的に生成します:
-
Cross Join のビルド側の統計上の行数が 1 であり、そのディストリビューションが Replicated である。
-
プローブ側のフィルター述語に、
>、<、>=、<=、またはBETWEENのようなビルド側の式との比較条件が含まれている。
ランタイムフィルターの種類
ランタイムフィルターは、以下の 2 つのディメンションで分類できます。
シャッフルスコープ別 (Hash Join の場合)
|
タイプ |
サポートバージョン |
シナリオ |
|
ローカル |
V2.0 以降 |
プローブ側のデータをシャッフルする必要がない場合に使用されます。ビルド側とプローブ側の結合キーが同じディストリビューションを共有している場合、ビルド側のデータがプローブ側にブロードキャストされる場合、またはビルド側のデータがプローブ側のディストリビューションに合わせてシャッフルされる場合に、ローカルランタイムフィルターを使用できます。このタイプは、スキャンされるデータ量と Hash Join によって処理されるデータ量のみを削減します。 |
|
グローバル |
V2.2 以降 |
プローブ側のデータをシャッフルする必要がある場合に使用されます。ランタイムフィルターはデータシャッフルの前に適用され、ネットワークトラフィックを削減します。 |
タイプを指定する必要はありません。エンジンが適応的に選択します。
フィルタータイプ別
|
タイプ |
サポートバージョン |
説明 |
|
ブルームフィルター |
V2.0 以降 |
このフィルターは確率的であり、誤検知が発生する可能性があります。つまり、一部のデータがフィルターされない場合があります。しかし、適用範囲が広く、ビルド側に大量のデータが含まれている場合でも高いフィルタリング効率を維持します。 |
|
In フィルター |
V2.0 以降 |
ビルド側の NDV (個別値の数) が小さい場合に使用されます。ビルド側のデータから HashSet を構築し、それをプローブ側に送信してフィルタリングします。必要なすべてのデータを正確にフィルタリングし、ビットマップ索引と組み合わせて使用できます。 |
|
MinMax フィルター |
V2.0 以降 |
ビルド側データの最小値と最大値をプローブ側に送信してフィルタリングします。メタデータを活用してファイル全体やデータのバッチをスキップし、I/O コストを削減できます。 |
|
ScalarFilter |
V4.2 以降 |
Cross Join のシナリオ専用に設計されています。ビルド側が正確に 1 行の場合、エンジンはスカラー値をプローブ側の ScanNode にプッシュダウンし、スキャンフェーズでフィルタリングを行います。 |
フィルタータイプを指定する必要はありません。Hologres はランタイムの結合条件に基づいて適応的にタイプを選択します。
Cross Join のランタイムフィルターサポート (V4.2 の新機能)
動機
ユーザーはしばしばスカラーサブクエリで集計値を計算し、その結果を大規模テーブルのフィルター条件として使用します。V4.2 以前は、この SQL パターンは実行計画において以下の問題がありました:
-
プローブ側の ScanNode は、ビルド側の値を早期フィルタリングに活用できず、全表スキャンを実行する必要がありました。
-
スキャンの後、Cross Join は行ごとの比較を行い、これが大きな I/O の無駄を引き起こしていました。
-
既存のランタイムフィルターは Hash Join のみを対象としており、Cross Join をサポートしていませんでした。
V4.2 では、この問題を解決するために ScalarFilter が導入されました。この機能はデフォルトで有効になっており、ユーザーによる操作は不要です。
典型的な SQL と実行計画の比較
典型的な SQL:
-- t1 is a large table; the aggregated min/max result from t2 is exactly 1 row.
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1);
最適化前 (ランタイムフィルターなし):
Cross Join
-> Seq Scan on t1 -- Full table scan, no early filtering
-> Aggregate -- Build side min/max, 1-row result
-> Seq Scan on t2
プローブ側はすべてのデータをスキャンする必要があり、その後 Cross Join が行を 1 つずつ比較するため、大きな I/O の無駄が発生します。
最適化後 (V4.2、ScalarFilter あり):
Cross Join
Runtime Filter Build Expr: (min(t2.a)), (max(t2.a))
-> Seq Scan on t1 -- Receives ScalarFilter, filters non-matching rows and RowGroups during scan
Runtime Filter Target Expr: (t1.a >= ${1}) AND (t1.a <= ${2})
-> Aggregate
-> Seq Scan on t2
ビルド側が評価された後、実際の値がプローブ側の ScanNode にプッシュダウンされ、フィルタリングがスキャンフェーズで完了できるようになります。
その他の典型的なシナリオ
-- Scenario 1: Single table + scalar subquery
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1);
-- Scenario 2: Multi-table + scalar subquery
SELECT * FROM t1, t3
WHERE (SELECT min(a) FROM t2) <= t1.a
AND t3.a <= (SELECT max(a) FROM t2);
-- Scenario 3: CTE + Cross Join
WITH r AS (SELECT MIN(a) AS lo, MAX(a) AS hi FROM t2)
SELECT t1.* FROM t1, r
WHERE t1.a >= r.lo AND t1.a <= r.hi;
仕組み
-
オプティマイザーによる識別:Cross Join のビルド側の統計上の行数が 1 で、ディストリビューションが Replicated の場合、オプティマイザーは述語から比較式を抽出し (
BETWEENを>=と<=に展開)、ScalarFilter の候補を生成します。 -
実行時のプッシュダウン:Open フェーズで、Cross Join はビルド側のデータを消費し、
build_exprを評価してスカラー値を生成し、ScalarFilter を構築してプローブ側の ScanNode に発行します。 -
ScanNode での早期フィルタリング:プローブ側の NiagaraScan が ScalarFilter を受け取った後、
target_exprのプレースホルダーを実際の値に置き換え、2 段階のフィルタリングを生成します:-
行レベルのフィルタリング (conjunct eval)。
-
RowGroup レベルのフィルタリング (ストレージエンジンが一致しないデータブロックをスキップします)。
-
制限事項
-
ビルド側は正確に 1 行でなければなりません。オプティマイザーは、ビルド側の統計上の行数が 1 で、ディストリビューションが Replicated の場合にのみ ScalarFilter を生成します。ランタイムでビルド側が正確に 1 行を生成しない場合、Hologres は
RT_CHECKエラーを報告します。 -
ビルド側の値が NULL の場合、Hologres は
FILTER_ALLフィルター (すべての行をフィルターする) を発行し、ビルド側のデータをクリアして Cross Join をショートサーキットさせ、空の結果セットを返します。 -
辞書エンコーディング列はサポートされていません。
target_exprによって参照されるプローブ側の列が辞書エンコーディング列 (DICTIONARY 型) の場合、ScalarFilter は適用されません。 -
サポートされているデータの型には、
INT8、UINT8、INT16、UINT16、INT32、UINT32、INT64、UINT64、DATE32、TIMESTAMP、FLOAT、DOUBLE、およびSTRINGが含まれます。サポートされていないデータの型の場合、ScalarFilter はプッシュダウンされず、エラーも報告されません。 -
サポートされている比較演算子には、
>、<、>=、および<=が含まれます。BETWEENは 2 つの範囲比較に展開されます。
GUC パラメーター
|
パラメーター |
タイプ |
デフォルト |
レベル |
説明 |
|
|
bool |
|
|
Cross Join に対して ScalarFilter を生成するかどうかを制御します。この機能はデフォルトで有効になっており、ユーザーによる操作は不要です。 |
この機能を無効にするには、次のコマンドを実行します:
SET hg_experimental_generate_runtime_scalar_filter = off;
EXPLAIN 出力の新しいフィールド
ScalarFilter が有効な場合、Cross Join ノードには次の情報が表示されます:
-
Runtime Filter Build Expr:ビルド側の式。例:(min(t2.a))および(max(t2.a))。 -
Runtime Filter Target Expr:プローブ側のターゲットフィルター式。例:(t1.a >= ${1}) AND (t1.a <= ${2})。この式では、${N}はプレースホルダーであり、ランタイムでfilter_idに対応するビルド側の値に置き換えられます。
パフォーマンス上の利点
TPC-DS 10 TB データセットでは、次の SQL パターンが一般的です:
DELETE FROM inventory
WHERE inv_date_sk >= (SELECT min(d_date_sk) FROM date_dim WHERE d_date BETWEEN 'INV_S_1' AND 'INV_E_1')
AND inv_date_sk <= (SELECT max(d_date_sk) FROM date_dim WHERE d_date BETWEEN 'INV_S_1' AND 'INV_E_1');
V4.2 の ScalarFilter 最適化により、クエリ時間は 2.672 秒から 0.552 秒に短縮され、パフォーマンスが約 4.8 倍向上します。
ランタイムフィルターの検証
以下の例は、さまざまなシナリオでランタイムフィルターの効果を検証する方法を示しています。
例 1:単一列の結合条件 (ローカルタイプ)
BEGIN;
CREATE TABLE test1 (x int, y int);
CALL set_table_property('test1', 'distribution_key', 'x');
CREATE TABLE test2 (x int, y int);
CALL set_table_property('test2', 'distribution_key', 'x');
END;
INSERT INTO test1 SELECT t, t FROM generate_series(1, 100000) t;
INSERT INTO test2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;
実行計画:
QUERY PLAN
Gather (cost=0.00..10.20 rows=1000 width=16)
[40:1 id=100002 dop=1 time=9/9/9ms rows=1000(1000/1000/1000) mem=16/16/16KB open=2/2/2ms get_next=7/7/7ms]
-> Hash Join (cost=0.00..10.16 rows=1000 width=16)
Hash Cond: (test1.x = test2.x)
Runtime Filter Cond: (test1.x = test2.x)
[id=8 dop=40 time=9/4/2ms rows=1000(39/25/17) mem=6/5/5KB open=4/1/0ms get_next=6/2/0ms]
-> Local Gather (cost=0.00..5.11 rows=1000000 width=8)
[id=3 dop=40 time=6/1/0ms rows=1000(39/25/17) mem=600/600/600B open=1/0/0ms get_next=6/1/0ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..5.10 rows=1000000 width=8)
Runtime Filter Target Expr: test1.x
[id=2 split_count=40 time=11/7/7ms rows=1000(39/25/17) mem=41/41/41KB open=11/7/7ms get_next=0/0/0ms scan_rows=1000000 25270/25000/24697))]
-> Hash (cost=5.00..5.00 rows=1000 width=8)
[id=7 dop=40 time=3/1/0ms rows=1000(39/25/17) mem=396/396/396KB open=3/1/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
-> Local Gather (cost=0.00..5.00 rows=1000 width=8)
[id=5 dop=40 time=1/0/0ms rows=1000(39/25/17) mem=0/0/0B open=1/0/0ms get_next=1/0/0ms local_dop=0/0/0]
-> Seq Scan on test2 (cost=0.00..5.00 rows=1000 width=8)
[id=4 split_count=40 time=1/0/0ms rows=1000(39/25/17) mem=528/528/528B open=1/0/0ms get_next=1/0/0ms scan_rows=1000(39/25/17)]
-
test2テーブルには 1,000 行、test1テーブルには 100,000 行あります。ビルド側とプローブ側のデータサイズの比率は 0.01 (0.1 未満) であり、ランタイムフィルターのデフォルトのトリガー条件を満たしています。 -
test1のプローブ側スキャンにはRuntime Filter Target Exprが表示されており、ランタイムフィルターがプッシュダウンされたことを示しています。 -
プローブ側では、
scan_rowsは 100,000 (ストレージから読み取られた行数) ですが、rowsは 1,000 (フィルタリング後の行数) です。この差がフィルタリング効果を示しています。
例 2:複数列の結合条件 (V2.1 以降、ローカルタイプ)
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x int, y int);
CREATE TABLE test2 (x int, y int);
END;
INSERT INTO test1 SELECT t, t FROM generate_series(1, 1000000) t;
INSERT INTO test2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x AND test1.y = test2.y;
実行計画:
QUERY PLAN
Gather (cost=0.00..10.46 rows=1000 width=16)
[40:1 id=100003 dop=1 time=6/6/6ms rows=1000(1000/1000/1000) mem=600/600/600B open=0/0/0ms get_next=6/6/6ms]
-> Hash Join (cost=0.00..10.43 rows=1000 width=16)
Hash Cond: ((test1.x = test2.x) AND (test1.y = test2.y))
Runtime Filter Cond: ((test1.x = test2.x) AND (test1.y = test2.y))
[id=8 dop=40 time=5/3/3ms rows=1000(1000/25/0) mem=40/2/1KB open=1/0/0ms get_next=4/3/3ms]
-> Local Gather (cost=0.00..5.11 rows=1000000 width=8)
[id=5 dop=40 time=4/3/3ms rows=1000(1000/25/0) mem=600/600/600B open=0/0/0ms get_next=4/3/3ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..5.10 rows=1000000 width=8)
Runtime Filter Target Expr: (test1.x AND test1.y)
[id=4 split_count=40 time=7/5/5ms rows=1000(1000/25/0) mem=49/11/9KB open=7/5/5ms get_next=1/0/0ms scan_rows=1000000(32768/25000/24576)]
-> Hash (cost=5.02..5.02 rows=40000 width=8)
[id=7 dop=40 time=1/0/0ms rows=40000(1000/1000/1000) mem=417/417/417KB open=1/0/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
-> Broadcast (cost=0.00..5.02 rows=40000 width=8)
[40:40 id=100002 dop=40 time=1/0/0ms rows=40000(1000/1000/1000) mem=0/0/0B open=1/0/0ms get_next=0/0/0ms * ]
-> Local Gather (cost=0.00..5.00 rows=1000 width=8)
[id=3 dop=40 time=1/0/0ms rows=1000(1000/25/0) mem=600/600/600B open=1/0/0ms get_next=0/0/0ms local_dop=1/1/1]
-> Seq Scan on test2 (cost=0.00..5.00 rows=1000 width=8)
[id=2 split_count=40 time=1/0/0ms rows=1000(1000/25/0) mem=528/528/528B open=0/0/0ms get_next=1/0/0ms scan_rows=1000(1000/1000/1000)]
-
結合条件には複数の列が含まれており、ランタイムフィルターも複数の列に対して生成されます。
-
ビルド側のデータはブロードキャストされるため、ローカルランタイムフィルターが使用されます。
例 3:グローバルタイプ (V2.2 以降、シャッフル結合)
SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x int, y int);
CREATE TABLE test2 (x int, y int);
END;
INSERT INTO test1 SELECT t, t FROM generate_series(1, 100000) t;
INSERT INTO test2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;
実行計画:
QUERY PLAN
-> Hash Join (cost=0.00..10.08 rows=1000 width=16)
Hash Cond: (test1.x = test2.x)
Runtime Filter Cond: (test1.x = test2.x)
[id=9 dop=40 time=10/8/8ms rows=1000(34/25/13) mem=6/6/5KB open=2/1/1ms get_next=8/7/7ms]
-> Redistribution (cost=0.00..5.07 rows=100000 width=8)
Hash Key: test1.x
[40:40 id=100002 dop=40 time=8/7/7ms rows=1289(46/32/20) mem=512/432/0B open=0/0/0ms get_next=8/7/7ms * ]
-> Local Gather (cost=0.00..5.01 rows=100000 width=8)
[id=3 dop=40 time=9/2/0ms rows=1289(1042/32/0) mem=600/600/600B open=0/0/0ms get_next=9/2/0ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..5.01 rows=100000 width=8)
Runtime Filter Target Expr: test1.x
[id=2 split_count=40 time=11/3/0ms rows=1289(1042/32/0) mem=50512/11849/528B open=11/3/0ms get_next=1/0/0ms scan_rows=100000(8192/7692/1696)]
-> Hash (cost=5.00..5.00 rows=1000 width=8)
[id=8 dop=40 time=2/1/1ms rows=1000(34/25/13) mem=396/396/396KB open=2/1/1ms get_next=0/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
-> Redistribution (cost=0.00..5.00 rows=1000 width=8)
Hash Key: test2.x
[40:40 id=100003 dop=40 time=2/1/1ms rows=1000(34/25/13) mem=0/0/0B open=0/0/0ms get_next=2/1/1ms * ]
-> Local Gather (cost=0.00..5.00 rows=1000 width=8)
[id=5 dop=40 time=1/0/0ms rows=1000(1000/25/0) mem=600/600/600B open=1/0/0ms get_next=1/0/0ms local_dop=1/1/1]
-> Seq Scan on test2 (cost=0.00..5.00 rows=1000 width=8)
[id=4 split_count=40 time=0/0/0ms rows=1000(1000/25/0) mem=528/528/528B open=0/0/0ms get_next=0/0/0ms scan_rows=1000(1000/1000/1000)]
プローブ側のデータは Hash Join オペレーターにシャッフルされます。エンジンはクエリを高速化するために自動的にグローバルランタイムフィルターを使用します。
例 4:ビットマップ索引を使用した In フィルター (V2.2 以降)
SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x text, y text);
CALL set_table_property('test1', 'distribution_key', 'x');
CALL set_table_property('test1', 'bitmap_columns', 'x');
CALL set_table_property('test1', 'dictionary_encoding_columns', '');
CREATE TABLE test2 (x text, y text);
CALL set_table_property('test2', 'distribution_key', 'x');
END;
INSERT INTO test1 SELECT t::text, t::text FROM generate_series(1, 10000000) t;
INSERT INTO test2 SELECT t::text, t::text FROM generate_series(1, 50) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;
実行計画:
QUERY PLAN
Gather (cost=0.00..11.70 rows=50 width=14)
[40:1 id=100002 dop=1 time=16/16/16ms rows=50(50/50/50) mem=2/2/2KB open=0/0ms get_next=16/16/16ms]
-> Hash Join (cost=0.00..11.70 rows=50 width=14)
Hash Cond: (test1.x = test2.x)
Runtime Filter Cond: (test1.x = test2.x)
[id=7 dop=40 time=15/9/3ms rows=50(3/1/0) mem=5132/3774/264B open=1/0/0ms get_next=14/8/3ms]
-> Local Gather (cost=0.00..6.26 rows=10000000 width=12)
[id=3 dop=40 time=14/8/3ms rows=50(3/1/0) mem=600/600/600B open=1/0/0ms get_next=14/8/3ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..6.06 rows=10000000 width=12)
Runtime Filter Target Expr: test1.x
[id=2 split_count=40 time=16/10/5ms rows=50(3/1/0) mem=67544/48945/528B open=16/9/5ms get_next=1/0/0ms scan_rows=7247692(250875/249920/248982) bitmap_used=50]
-> Hash (cost=5.00..5.00 rows=50 width=2)
[id=6 dop=40 time=1/0/0ms rows=61(3/1/1) mem=534/530/521KB open=1/0/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=512/512/512KB]
-> Local Gather (cost=0.00..5.00 rows=50 width=2)
[id=5 dop=40 time=1/0/0ms rows=50(3/1/0) mem=600/600/600B open=0/0/0ms get_next=1/0/0ms local_dop=1/1/1]
-> Seq Scan on test2 (cost=0.00..5.00 rows=50 width=2)
[id=4 split_count=40 time=1/0/0ms rows=50(3/1/0) mem=528/528/528B open=1/0/0ms get_next=0/0/0ms scan_rows=50(3/1/1)]
プローブ側のスキャンオペレーターはビットマップ索引を使用します。In フィルターは正確なフィルタリングを提供し、50 行のみを残します。スキャンオペレーターの scan_rows の値は 700 万を超えており、元の 1,000 万行よりも少なくなっています。これは、In フィルターがストレージエンジンにプッシュダウンされ、I/O オーバーヘッドを削減できるためです。In フィルターとビットマップ索引を組み合わせることで、結合キーが STRING 型の場合に大きな効果が得られます。
例 5:I/O 削減のための MinMax フィルター (V2.2 以降)
SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS test1, test2;
BEGIN;
CREATE TABLE test1 (x int, y int);
CALL set_table_property('test1', 'distribution_key', 'x');
CREATE TABLE test2 (x int, y int);
CALL set_table_property('test2', 'distribution_key', 'x');
END;
INSERT INTO test1 SELECT t::int, t::int FROM generate_series(1, 10000000) t;
INSERT INTO test2 SELECT t::int, t::int FROM generate_series(1, 100000) t;
ANALYZE test1;
ANALYZE test2;
EXPLAIN ANALYZE SELECT * FROM test1 JOIN test2 ON test1.x = test2.x;
実行計画:
QUERY PLAN
Gather (cost=0.00..15.68 rows=100000 width=16)
[40:1 id=100002 dop=1 time=5/5/5ms rows=100000(100000/100000/100000) mem=600/600/600B open=0/0/0ms get_next=5/5/5ms]
-> Hash Join (cost=0.00..11.98 rows=100000 width=16)
Hash Cond: (test1.x = test2.x)
Runtime Filter Cond: (test1.x = test2.x)
[id=7 dop=40 time=5/4/4ms rows=100000(2639/2500/2406) mem=97/92/89KB open=1/0/0ms get_next=4/3/3ms]
-> Local Gather (cost=0.00..6.14 rows=10000000 width=8)
[id=3 dop=40 time=5/3/3ms rows=100000(2639/2500/2406) mem=600/600/600B open=1/0/0ms get_next=4/3/3ms local_dop=1/1/1]
-> Seq Scan on test1 (cost=0.00..6.00 rows=10000000 width=8)
Runtime Filter Target Expr: test1.x
[id=2 split_count=40 time=6/6/5ms rows=100000(2639/2500/2406) mem=61/60/59KB open=6/5/5ms get_next=0/0/0ms scan_rows=327680(8192/8192/8192)]
-> Hash (cost=5.01..5.01 rows=100000 width=8)
[id=6 dop=40 time=1/0/0ms rows=100000(2639/2500/2406) mem=463/460/458KB open=1/0/0ms get_next=1/0/0ms rehash=1/1/1 hash_mem=384/384/384KB]
-> Local Gather (cost=0.00..5.01 rows=100000 width=8)
[id=5 dop=40 time=1/0/0ms rows=100000(2639/2500/2406) mem=600/600/600B open=0/0/0ms get_next=1/0/0ms local_dop=1/1/1]
-> Seq Scan on test2 (cost=0.00..5.01 rows=100000 width=8)
[id=4 split_count=40 time=1/0/0ms rows=100000(2639/2500/2406) mem=528/528/528B open=1/0/0ms get_next=0/0/0ms scan_rows=100000(2639/2500/2406)]
プローブ側のスキャンオペレーターは、ストレージエンジンから 32 万行強しか読み取っておらず、元の 1,000 万行よりもはるかに少なくなっています。これは、ランタイムフィルターがストレージエンジンにプッシュダウンされ、データバッチのメタデータを使用してバッチ全体を一度にフィルタリングするためです。これにより、I/O オーバーヘッドを大幅に削減できます。このフィルタータイプは、結合キーが数値型で、ビルド側の値の範囲がプローブ側の範囲よりも小さい場合に最も効果的です。
例 6:TopN ランタイムフィルター (V4.0 以降)
SQL ステートメントに topN オペレーターが含まれている場合、Hologres はすべての結果を計算するのではなく、動的なフィルターを生成して早期にデータをプルーニングします。
SELECT o_orderkey FROM orders ORDER BY o_orderdate LIMIT 5;
実行計画:
QUERY PLAN
Limit (cost=0.00..116554.70 rows=0 width=8)
-> Sort (cost=0.00..116554.70 rows=100 width=12)
Sort Key: o_orderdate
[id=6 dop=1 time=317/317/317ms rows=5(5/5/5) mem=1/1/1KB open=317/317/317ms get_next=0/0/0ms]
-> Gather (cost=0.00..116554.25 rows=100 width=12)
[20:1 id=100002 dop=1 time=317/317/317ms rows=100(100/100/100) mem=6/6/6KB open=0/0/0ms get_next=317/317/317ms * ]
-> Limit (cost=0.00..116554.25 rows=0 width=12)
-> Sort (cost=0.00..116554.25 rows=150000000 width=12)
Sort Key: o_orderdate
Runtime Filter Sort Column: o_orderdate
[id=3 dop=20 time=318/282/258ms rows=100(5/5/5) mem=96/96/96KB open=318/282/258ms get_next=1/0/0ms]
-> Local Gather (cost=0.00..9.59 rows=150000000 width=12)
[id=2 dop=20 time=316/280/256ms rows=1372205(68691/68610/68498) mem=0/0/0B open=0/0/0ms get_next=316/280/256ms local_dop=1/1/1 * ]
-> Seq Scan on orders (cost=0.00..8.24 rows=150000000 width=12)
Runtime Filter Target Expr: o_orderdate
[id=1 split_count=20 time=286/249/222ms rows=1372205(68691/68610/68498) mem=179/179/179KB open=0/0/0ms get_next=286/249/222ms physical_reads=27074(1426/1353/1294) scan_rows=144867963(7324934/7243398/7172304)]
Query id:[1001003033996040311]
QE version: 2.0
Query Queue: init_warehouse.default_queue
======================cost======================
Total cost:[343] ms
Optimizer cost:[13] ms
Build execution plan cost:[0] ms
Init execution plan cost:[6] ms
Start query cost:[6] ms
- Queue cost: [0] ms
- Wait schema cost:[0] ms
- Lock query cost:[0] ms
- Create dataset reader cost:[0] ms
- Create split reader cost:[0] ms
Get result cost:[318] ms
- Get the first block cost:[318] ms
====================resource====================
Memory: total 7 MB. Worker stats: max 3 MB, avg 3 MB, min 3 MB, max memory worker id: 189*****.
CPU time: total 5167 ms. Worker stats: max 2610 ms, avg 2583 ms, min 2557 ms, max CPU time worker id: 189*****.
DAG CPU time stats: max 5165 ms, avg 2582 ms, min 0 ms, cnt 2, max CPU time dag id: 1.
Fragment CPU time stats: max 5137 ms, avg 1721 ms, min 0 ms, cnt 3, max CPU time fragment id: 2.
Ec wait time: total 90 ms. Worker stats: max 46 ms, max(max) 2 ms, avg 45 ms, min 44 ms, max ec wait time worker id: 189*****, max(max) ec wait time worker id: 189*****.
Physical read bytes: total 799 MB. Worker stats: max 400 MB, avg 399 MB, min 399 MB, max physical read bytes worker id: 189*****.
Read bytes: total 898 MB. Worker stats: max 450 MB, avg 449 MB, min 448 MB, max read bytes worker id: 189*****.
DAG instance count: total 3. Worker stats: max 2, avg 1, min 1, max DAG instance count worker id: 189*****.
Fragment instance count: total 41. Worker stats: max 21, avg 20, min 20, max fragment instance count worker id: 189*****.
TopN ランタイムフィルターがない場合、ScanNode は orders テーブルからすべてのデータブロックを読み取り、それを TopN ノードに渡します。TopN ノードはヒープソートを使用して、これまでに見つかった上位 5 行を維持します。
例えば、各データブロックには約 8,192 行が含まれています。最初のブロックが処理された後、TopN はそのブロック内で 5 番目にランク付けされた o_orderdate を知ります。それが 1995-01-01 だとします。Scan ノードが 2 番目のブロックを読み取る際、1995-01-01 をフィルター条件として使用し、o_orderdate <= 1995-01-01 の行のみを TopN に送信します。しきい値は動的に更新されます。2 番目のブロックで 5 番目にランク付けされた o_orderdate が小さい場合、TopN は古いしきい値を新しい値に置き換えます。
EXPLAIN 出力を使用して、オプティマイザーによって生成された TopN ランタイムフィルターを表示できます:
-> Limit (cost=0.00..116554.25 rows=0 width=12)
-> Sort (cost=0.00..116554.25 rows=150000000 width=12)
Sort Key: o_orderdate
Runtime Filter Sort Column: o_orderdate
[id=3 dop=20 time=318/282/258ms rows=100(5/5/5) mem=96/96/96KB open=318/282/258ms get_next=1/0/0ms]
TopN ノードに Runtime Filter Sort Column が存在することは、このノードが TopN ランタイムフィルターを生成することを示しています。
例 7:ScalarFilter を使用した Cross Join (V4.2 以降)
SET hg_experimental_enable_result_cache = OFF;
DROP TABLE IF EXISTS t1, t2;
BEGIN;
CREATE TABLE t1 (a int, b int);
CREATE TABLE t2 (a int, b int);
END;
INSERT INTO t1 SELECT t, t FROM generate_series(1, 1000000) t;
INSERT INTO t2 SELECT t, t FROM generate_series(1, 1000) t;
ANALYZE t1;
ANALYZE t2;
EXPLAIN ANALYZE
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1)
AND a <= (SELECT max(a) FROM t2 WHERE b BETWEEN 0 AND 1);
実行計画のハイライト:
-
Cross Join ノードには
Runtime Filter Build Expr: (min(t2.a)), (max(t2.a))が表示されます。 -
t1 の ScanNode には
Runtime Filter Target Expr: (t1.a >= ${1}) AND (t1.a <= ${2})が表示されます。 -
t1 の場合、
scan_rowsはストレージから読み取られた行数を反映しますが、rowsは ScalarFilter が適用された後に大幅に減少します。この差がフィルターの効果を検証します。
比較のために機能を無効化:
SET hg_experimental_generate_runtime_scalar_filter = off;
EXPLAIN ANALYZE
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1)
AND a <= (SELECT max(a) FROM t2 WHERE b BETWEEN 0 AND 1);
機能が無効になると、実行計画には Runtime Filter Build Expr と Target Expr フィールドが表示されなくなり、t1 のスキャンは全表スキャンに戻ります。
バージョン履歴
|
バージョン |
新機能 |
|
V2.0 |
Hash Join でのランタイムフィルターのサポート (ローカルタイプ、ブルーム、In、MinMax フィルターを含む)。 |
|
V2.1 |
複数列の結合条件を持つランタイムフィルターのサポート。 |
|
V2.2 |
グローバルランタイムフィルターのサポート (シャッフル結合用);In フィルターはビットマップ索引と組み合わせ可能。 |
|
V4.0 |
TopN ランタイムフィルターのサポート。 |
|
V4.2 |
Cross Join ランタイムフィルター (ScalarFilter) のサポート。スカラーサブクエリを使用して大規模テーブルをフィルタリングするシナリオを最適化。 |