全部產品
Search
文件中心

Hologres:Runtime Filter 多表 Join 加速

更新時間:Jun 09, 2026

在多表 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;

工作原理

  1. 最佳化器識別:當 Cross Join 的 build 端統計行數為 1 且分布為 Replicated 時,從謂詞中提取比較運算式(BETWEEN 展開為 >= + <=),產生 ScalarFilter 候選。

  2. 執行期下推:Cross Join 在 Open 階段消費 build 端資料後,對 build_expr 求值得到標量值,構建 ScalarFilter 發布給 probe 端 ScanNode。

  3. 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 參數

參數名

類型

預設值

層級

說明

hg_experimental_generate_runtime_scalar_filter

bool

true

PGC_USERSET

是否為 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),最佳化「標量子查詢 + 大表過濾」情境