在多表 Join 查詢中,未參與 Join 匹配的資料會增加 I/O 開銷並降低查詢效能。Runtime Filter 在多表 Join 情境下自動產生輕量過濾器,在資料掃描階段提前裁剪無效資料,減少 I/O 開銷並提升查詢效能。該功能自 Hologres V2.0 開始支援 Hash Join,V4.2 擴充至 Cross Join 情境。
背景資訊
應用情境
Hologres 從 V2.0 版本開始支援 Runtime Filter,通常應用在多表(兩表及以上)Join 的 Hash Join 情境,尤其是大表 Join 小表的情境中,無需手動設定,最佳化器和執行引擎會在查詢時自動最佳化 Join 過程的過濾行為,從而降低 I/O 開銷,提升 Join 的查詢效能。
從 V4.2 版本開始,Runtime Filter 的能力進一步擴充到 Cross Join 情境,專門最佳化「標量子查詢彙總結果作為大表過濾條件」這類高頻 SQL 模式(詳見Cross Join 支援 Runtime Filter(V4.2 新增))。
原理介紹
Hash Join 情境下的 Runtime Filter
兩個表 Join 時,會通過一張表構建 Hash 表,然後匹配另一張表的資料。Join 過程涉及兩個端:
-
build 端:構建 Hash 表的一側,對應執行計畫中的 Hash-節點。
-
probe 端:讀取資料並與 build 端的 Hash 表進行匹配的一側。
通常小表作為 build 端,大表作為 probe 端。
Runtime Filter 的原理是:利用 build 端的資料分布產生輕量過濾器,發送給 probe 端對資料進行裁剪,減少 probe 端參與 Hash Join 及網路傳輸的資料量,從而提升 Join 效能。 因此該功能更適用於大小表 Join 且表資料量相差較大的情境,效能將會比普通 Join 有更多提升。
Cross Join 情境下的 Runtime Filter(V4.2 新增)
V4.2 之前,Runtime Filter 僅覆蓋 Hash Join。但使用者經常使用標量子查詢計算彙總值(如 min / max),並將結果作為大表的過濾條件,這類 SQL 在執行計畫中產生的是 Cross Join:
-
build 端:標量子查詢結果,恰好 1 行。
-
probe 端:大表掃描。
V4.2 引入了新的 Filter 類型 —— ScalarFilter,由最佳化器自動識別並將 build 端求值後的標量值下推給 probe 端 ScanNode,掃描階段即完成行級與 RowGroup 級過濾,避免大表全表掃描後再逐行比較。
使用限制和觸發條件
使用限制
-
僅 Hologres V2.0 及以上版本支援 Runtime Filter。
-
Hash Join 情境下,V2.0 版本僅支援 Join 條件中只有一個欄位。從 V2.1 版本開始,支援多個欄位的 Runtime Filter。
-
僅 V4.0 及以上版本支援 TopN Runtime Filter,用以提升單表 TopN 計算情境的效能。
-
僅 V4.2 及以上版本支援 Cross Join Runtime Filter(ScalarFilter)。
觸發條件
Hash Join 情境
Runtime Filter 由引擎自動觸發,需同時滿足以下條件:
-
probe 端資料量在 100,000 行及以上。
-
掃描資料量比例:build 端 / probe 端 ≤ 0.1(比例越小,越容易觸發)。
-
Join 輸出資料量比例:build 端 / probe 端 ≤ 0.1(比例越小,越容易觸發)。
Cross Join 情境(V4.2+)
由最佳化器自動產生 ScalarFilter,需滿足:
-
Cross Join 的 build 端統計行數為 1,且分布為 Replicated。
-
probe 端的過濾謂詞中包含 build 端運算式的比較條件(
>、<、>=、<=或BETWEEN)。
Runtime Filter 的類型
Runtime Filter 可從以下兩個維度進行分類。
按 Shuffle 維度劃分(適用於 Hash Join)
|
類型 |
支援版本 |
適用情境 |
|
Local |
V2.0+ |
probe 端資料不需要 Shuffle 時使用。build 端和 probe 端的 Join Key 為同一種分布方式,或 build 端資料 broadcast 給 probe 端,或 build 端資料按照 probe 端的分布方式 Shuffle 給 probe 端,均可使用 Local 類型。僅減少資料掃描量及參與 Hash Join 計算的資料量。 |
|
Global |
V2.2+ |
probe 端資料需要 Shuffle 時使用。Runtime Filter 在資料 Shuffle 之前做過濾,可減少資料的網路傳輸量。 |
類型無需手動指定,引擎會自適應選擇。
按 Filter 類型劃分
|
類型 |
支援版本 |
說明 |
|
Bloom Filter |
V2.0+ |
具有一定假陽性,可能少過濾部分資料,但應用範圍廣,在 build 端資料量較多時仍有較高的過濾效率。 |
|
In Filter |
V2.0+ |
build 端資料 NDV(非重複值個數)較小時使用。使用 build 端資料構建 HashSet 發送給 probe 端過濾,可過濾所有應過濾的資料,且能與 Bitmap 索引結合使用。 |
|
MinMax Filter |
V2.0+ |
根據 build 端資料的最大值和最小值發送給 probe 端做過濾。可根據中繼資料資訊直接過濾掉檔案或一個 Batch 的資料,減少 I/O 成本。 |
|
ScalarFilter |
V4.2+ |
專用於 Cross Join 情境。當 build 端恰好為 1 行時,引擎將標量值下推至 probe 端 ScanNode,在掃描階段直接過濾。 |
Filter 類型無需手動指定,Hologres 會根據運行時 Join 情況自適應使用。
Cross Join 支援 Runtime Filter(V4.2 新增)
需求背景
使用者經常使用標量子查詢計算彙總值,並將結果作為大表的過濾條件。這類 SQL 在 V4.2 之前的執行計畫中存在如下問題:
-
probe 端 ScanNode 無法利用 build 端的值提前過濾,必須全表掃描。
-
掃描完成後由 Cross Join 逐行比較,造成大量無效 I/O。
-
已有的 Runtime Filter 僅覆蓋 Hash Join,不支援 Cross Join。
V4.2 引入 ScalarFilter,從根本上解決該問題,使用者無需任何額外操作,功能預設開啟。
典型 SQL 與執行計畫對比
典型 SQL:
-- t1 大表,t2 彙總後 min/max 結果恰好 1 行
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1);
最佳化前(無 Runtime Filter):
Cross Join
-> Seq Scan on t1 -- 全表掃描,無法提前過濾
-> Aggregate -- build 端 min/max,結果 1 行
-> Seq Scan on t2
probe 端必須掃描全量資料,再由 Cross Join 逐行比較,存在大量無效 I/O。
最佳化後(V4.2 ScalarFilter):
Cross Join
Runtime Filter Build Expr: (min(t2.a)), (max(t2.a))
-> Seq Scan on t1 -- 收到 ScalarFilter,掃描時直接過濾不滿足條件的行和 RowGroup
Runtime Filter Target Expr: (t1.a >= ${1}) AND (t1.a <= ${2})
-> Aggregate
-> Seq Scan on t2
build 端求值後,將實際值下推給 probe 端 ScanNode,掃描階段即完成過濾。
更多典型情境
-- 情境一:單表 + 標量子查詢
SELECT * FROM t1
WHERE a >= (SELECT min(a) FROM t2 WHERE b BETWEEN 0 AND 1);
-- 情境二:多表 + 標量子查詢
SELECT * FROM t1, t3
WHERE (SELECT min(a) FROM t2) <= t1.a
AND t3.a <= (SELECT max(a) FROM t2);
-- 情境三: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 的 build 端統計行數為 1 且分布為 Replicated 時,從謂詞中提取比較運算式(
BETWEEN展開為>=+<=),產生 ScalarFilter 候選。 -
執行期下推:Cross Join 在 Open 階段消費 build 端資料後,對
build_expr求值得到標量值,構建 ScalarFilter 發布給 probe 端 ScanNode。 -
ScanNode 提前過濾:probe 端 NiagaraScan 接收到 ScalarFilter 後,用實際值替換
target_expr中的預留位置,同時產生兩級過濾:-
行級過濾(conjunct eval)。
-
RowGroup 級過濾(儲存引擎層直接跳過不滿足條件的資料區塊)。
-
使用限制
-
build 端必須恰好 1 行。最佳化器僅在 build 端統計行數為 1 且為 Replicated 分布時產生 ScalarFilter。如果運行時 build 端非 1 行,會觸發
RT_CHECK報錯。 -
build 端值為 NULL 時,發布
FILTER_ALL(過濾所有行),同時清空 build 端資料短路 Cross Join,直接返回空結果。 -
不支援字典列。如果 probe 端
target_expr引用的列是字典編碼列(DICTIONARY type),則不會應用 ScalarFilter。 -
支援的資料類型:
INT8/UINT8/INT16/UINT16/INT32/UINT32/INT64/UINT64/DATE32/TIMESTAMP/FLOAT/DOUBLE/STRING。不支援的類型會跳過 ScalarFilter(不報錯,僅不下推)。 -
支援的比較操作符:
>、<、>=、<=。BETWEEN會被展開為兩個範圍比較。
GUC 參數
|
參數名 |
類型 |
預設值 |
層級 |
說明 |
|
|
bool |
|
|
是否為 Cross Join 產生 ScalarFilter 類型的 Runtime Filter。預設開啟,使用者無需額外操作。 |
關閉該功能:
SET hg_experimental_generate_runtime_scalar_filter = off;
EXPLAIN 輸出新增欄位
啟用 ScalarFilter 後,Cross Join 節點會顯示:
-
Runtime Filter Build Expr:build 端運算式,例如(min(t2.a))、(max(t2.a))。 -
Runtime Filter Target Expr:probe 端目標過濾運算式,例如(t1.a >= ${1}) AND (t1.a <= ${2}),其中${N}為運行時替換的filter_id對應的 build 端值預留位置。
效能收益
在 TPC-DS 10TB 資料集中,存在如下典型 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 倍。
驗證 Runtime Filter
以下樣本協助您驗證 Runtime Filter 在不同情境下的效果。
樣本 1:Join 條件中只有 1 列(Local 類型)
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 行,build 端和 probe 端的資料量比例為 0.01(小於 0.1),滿足 Runtime Filter 的預設觸發條件。 -
probe 端的
test1表出現Runtime Filter Target Expr節點,表示 Runtime Filter 已下推。 -
probe 端
scan_rows為 100,000 行(從儲存中讀取的資料),rows為 1,000 行(過濾後的行數),可據此觀察過濾效果。
樣本 2:Join 條件中有多列(V2.1+,Local 類型)
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)]
-
Join 條件有多列,Runtime Filter 也產生了多列。
-
build 端資料採用 broadcast 方式,可使用 Local 類型的 Runtime Filter。
樣本 3:Global 類型(V2.2+,支援 Shuffle Join)
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)]
probe 端資料被 Shuffle 到 Hash Join 運算元,引擎自動使用 Global Runtime Filter 加速查詢。
樣本 4:In 類型結合 Bitmap 索引(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)]
probe 端的 scan 運算元使用了 bitmap。In Filter 可精確過濾,過濾後僅剩 50 行。scan 運算元中的 scan_rows 為 700 多萬(低於原始行數 1,000 萬),這是因為 In Filter 可下推到儲存引擎,減少 I/O 成本。In 類型 Runtime Filter 結合 bitmap 在 Join Key 為 STRING 類型時效果顯著。
樣本 5:MinMax 類型減少 I/O(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)]
probe 端 scan 運算元從儲存引擎讀取的行數為 32 萬多,遠低於原始行數 1,000 萬。這是因為 Runtime Filter 被下推到儲存引擎,利用一個 Batch 資料的 meta 資訊整批過濾,有可能大量減少 I/O 成本。該類型通常在 Join Key 為數實值型別且 build 端範圍範圍小於 probe 端時效果明顯。
樣本 6:TopN Runtime Filter(V4.0+)
當 SQL 陳述式包含 topN 運算元時,Hologres 不會計算所有結果,而是產生動態 Filter 提前對資料進行過濾。
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 Filter 時,Scan 節點讀取 orders 表的每個資料區塊並傳給 TopN 節點,TopN 節點用堆排序維護當前已見資料中排名前 5 的行。
例如:每個資料區塊約含 8,192 行。處理完第一個塊後,TopN 就知道該塊中第 5 名的 o_orderdate。假設它是 1995-01-01。Scan 節點在讀第二個塊時,就用 1995-01-01 作過濾條件,只發送 o_orderdate <= 1995-01-01 的行給 TopN。閾值會動態更新,如果第二個塊中第 5 名的 o_orderdate 更小,TopN 就用新值替換舊閾值。
通過 EXPLAIN 命令可查看最佳化器產生的 TopN Runtime Filter:
-> 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 Runtime Filter。
樣本 7:Cross Join 使用 ScalarFilter(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 的 Scan 節點輸出
Runtime Filter Target Expr: (t1.a >= ${1}) AND (t1.a <= ${2})。 -
t1 的
scan_rows仍為儲存中讀取的行數,經 ScalarFilter 過濾後rows顯著下降,可據此驗證過濾效果。
關閉功能驗證(用於迴歸對比):
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 的 Runtime Filter(Local 類型,Bloom / In / MinMax Filter) |
|
V2.1 |
支援多欄位 Join 的 Runtime Filter |
|
V2.2 |
支援 Global 類型 Runtime Filter(Shuffle Join);In Filter 可結合 Bitmap 索引 |
|
V4.0 |
支援 TopN Runtime Filter |
|
V4.2 |
支援 Cross Join Runtime Filter(ScalarFilter),最佳化「標量子查詢 + 大表過濾」情境 |