對於每條流式資料,可以關聯一個外部維表資料來源,為Realtime ComputeFlink版提供資料關聯查詢。
維表 JOIN 概述
維表 JOIN(Lookup Join)通過對每條流資料在處理時間查詢外部維表,將維表欄位補全到主流中。常用於事件流補充字典、維度資訊等情境。
維表 JOIN 的核心環節:
文法:使用
FOR SYSTEM_TIME AS OF PROCTIME()表示對每條流資料查詢維表的當前資料,而非維錶快照。Cache 策略:在連接器層緩衝維表資料,降低外部存取壓力。
Lookup 行為:通過 LOOKUP hint 配置同步/非同步、緩衝容量、重試策略等。
Join 物理策略:通過 SHUFFLE_HASH 等 hint 控制 Shuffle 行為,最佳化資料扭曲。
Job 級配置:通過 SET 命令調整
table.exec.*系列全域參數。
維表 JOIN 文法
SELECT column-names
FROM table1 [AS <alias1>]
[LEFT] JOIN table2 FOR SYSTEM_TIME AS OF PROCTIME() [AS <alias2>]
ON table1.column-name1 = table2.key-name1;必須加上
FOR SYSTEM_TIME AS OF PROCTIME(),表示 JOIN 維表當前時刻所看到的每條資料。ON 條件中必須包含維表實際能支援隨機尋找的欄位的等值條件。
ON 條件中支援對源表欄位使用 CAST 等類型轉換函式。如源表與維表欄位類型不一致,可在源表欄位上進行類型轉換以匹配維表欄位類型。
使用限制與注意事項
維表 JOIN 僅支援對當前時刻維錶快照的關聯。
維表支援 INNER JOIN 和 LEFT JOIN,不支援 RIGHT JOIN 或 FULL JOIN。
如有一對一 JOIN 需求,請確保串連條件中包含了維表中具有唯一性欄位的等值串連條件。
對每條流式資料,只會關聯當時維表的最新版本資料,即JOIN行為只發生在處理時間(Processing Time)。如果JOIN行為發生後,維表中的資料發生了變化(新增、更新或刪除),則已關聯的維表資料不會被同步變化。具體的維表的行為請參見對應連接器行為。
維表 Cache 策略
大部分連接器的維表 JOIN 支援 Cache 策略,不同連接器的支援情況略有差異,請查閱對應連接器文檔確認。通用 Cache 策略如下:
策略 | 行為 |
None(預設) | 無緩衝。 |
LRU | 緩衝維表裡的部分資料。源表的每條資料都會觸發系統先在 Cache 中尋找資料,未命中時再查物理維表。 |
ALL | 緩衝維表裡的所有資料。Job 運行前,系統會將維表中所有資料載入到 Cache 中,之後所有尋找都通過 Cache 完成;未命中視為 KEY 不存在。可通過連接器參數配置周期性或定時重新載入維表資料。適用於維表資料量較小、對即時性要求不高且需要極致吞吐的情境。 |
需根據業務需求在即時性和效能之間權衡。如對即時性要求極高,可不使用 Cache,直接從維表讀取。
使用 Cache 時,可配合 LRU 和 TTL 維持較新的快取資料。TTL 可設定較短(例如幾秒至幾十秒),定期從源表載入資料。
使用 ALL 緩衝策略時,注意節點記憶體大小,防止 OOM。
ALL 緩衝策略下系統非同步載入維表資料,需為維表 JOIN 節點增加記憶體,建議增加的記憶體大小為遠端資料表資料量的兩倍。
維表 JOIN 調優
調優參數分類與傳遞方式
維表 JOIN 調優涉及三類參數,每類參數的傳遞方式不同,用錯方式 hint 會被靜默忽略。
參數類別 | 典型參數 | 傳遞方式 | 範圍 |
連接器 WITH 選項 |
| 建表語句的 | 單張表 |
LOOKUP hint 選項 |
|
| 單個 Join 操作 |
Flink Job 級 TableConfig | 以 | 作業開頭使用 | 整個作業 |
OPTIONS()hint 僅用於覆蓋建表語句的 WITH 選項,不接受table.exec.*類 TableConfig;寫錯會被靜默忽略。LOOKUP()hint 只識別上表列出的固定 key,不接受 Flink 配置項全名;例如調整 async 緩衝,必須寫'capacity'='100',寫'table.exec.async-lookup.buffer-capacity'='100'不生效。如需調整 hint 範圍外的全域參數,請使用
SET 'xxx' = 'yyy';。
通過 OPTIONS hint 覆蓋連接器選項
OPTIONS() hint 用於在 SQL 中臨時覆蓋某張表的 WITH 選項,無需修改建表語句。常用於調整 Cache 行為。
SELECT t.id, t.name, w.phoneNumber
FROM kafka_input AS t
LEFT JOIN phoneNumber /*+ OPTIONS(
'lookup.cache' = 'PARTIAL',
'lookup.partial-cache.max-rows' = '1000'
) */ FOR SYSTEM_TIME AS OF PROCTIME() AS w
ON t.name = w.name;OPTIONS 中傳入的 key 必須為連接器本身支援的 WITH 選項。
上例中的 lookup.cache、lookup.partial-cache.max-rows 是基於 Flink 通用 LookupCache 介面(FLIP-221)的連接器選項,僅適用於實現了該介面的新版連接器(如 Fluss、JDBC 等)。前文「維表 Cache 策略」中描述的 None/LRU/ALL 策略對應歷史介面的連接器,兩者通過不同的 WITH 選項配置。請查閱對應連接器文檔確認其支援的 Cache 選項與策略,避免使用不匹配的參數導致配置不生效。
通過 LOOKUP hint 配置 Lookup 行為
LOOKUP hint 功能與社區保持一致,用於在單個 Join 操作上配置維表的同步、非同步、重試及 shuffle 策略。詳情參見 Apache Flink Lookup Hint。
僅 VVR 8.0 及以上版本支援 LOOKUP hint。
僅 VVR 8.0.8 及以上版本支援通過
'shuffle' = 'true'配置 shuffle 策略。VVR 8.0 以上支援使用別名;如果維表定義了別名,hint 中必須使用別名。
支援的選項
選項 | 含義 | 取值 |
| 指定 hint 作用的維表名或別名 | 表名/別名字串 |
| 是否啟用非同步 lookup |
|
| 非同步 lookup 的輸出順序 |
|
| 非同步 lookup 緩衝隊列容量 | 整數 |
| 非同步 lookup 逾時時間 | 時間長度(如 |
| 觸發重試的條件 | 當前僅支援 |
| 重試策略 | 當前僅支援 |
| 固定稍候再試 | 時間長度(如 |
| 最大重試次數 | 整數 |
| 是否在維表 Join 前做 shuffle |
|
shuffle 選項行為
shuffle 選項影響維表 Join 的 shuffle 策略,不同情境表現如下。
情境 | 聯結策略 |
不配置 | 使用引擎預設的 shuffle 策略。 |
不配置 | 使用引擎預設的 shuffle 策略。 |
配置 | 預設使用 SHUFFLE_HASH 策略,含義請參見《SHUFFLE_HASH》。 |
配置 | 使用表連接器的自訂 shuffle 策略。 |
目前僅流式資料湖倉Paimon會提供自訂shuffle策略,具體會在Join欄位包含全部分桶欄位的情況下基於bucket進行shuffle。
程式碼範例
-- 只對維表 dim1 配置維表聯結 shuffle 策略
SELECT /*+ LOOKUP('table'='dim1', 'shuffle' = 'true') */ ...
FROM src AS T
LEFT JOIN dim1 FOR SYSTEM_TIME AS OF PROCTIME() ON T.a = dim1.a
LEFT JOIN dim2 FOR SYSTEM_TIME AS OF PROCTIME() ON T.b = dim2.b;
-- 同時對維表 dim1、dim2 配置維表聯結 shuffle 策略
SELECT /*+ LOOKUP('table'='dim1', 'shuffle' = 'true'), LOOKUP('table'='dim2', 'shuffle' = 'true') */ ...
FROM src AS T
LEFT JOIN dim1 FOR SYSTEM_TIME AS OF PROCTIME() ON T.a = dim1.a
LEFT JOIN dim2 FOR SYSTEM_TIME AS OF PROCTIME() ON T.b = dim2.b;
-- 維表 dim1 使用別名 D1 時,hint 中必須使用別名
SELECT /*+ LOOKUP('table'='D1', 'shuffle' = 'true') */ ...
FROM src AS T
LEFT JOIN dim1 FOR SYSTEM_TIME AS OF PROCTIME() AS D1 ON T.a = D1.a
LEFT JOIN dim2 FOR SYSTEM_TIME AS OF PROCTIME() AS D2 ON T.b = D2.b;
-- 同時對維表 dim1、dim2 通過別名配置維表聯結 shuffle 策略
SELECT /*+ LOOKUP('table'='D1', 'shuffle' = 'true'), LOOKUP('table'='D2', 'shuffle' = 'true') */ ...
FROM src AS T
LEFT JOIN dim1 FOR SYSTEM_TIME AS OF PROCTIME() AS D1 ON T.a = D1.a
LEFT JOIN dim2 FOR SYSTEM_TIME AS OF PROCTIME() AS D2 ON T.b = D2.b;
通過 SET 配置 Job 級 TableConfig
table.exec.* 等 TableConfig 屬於 Job 級全域配置,必須通過 SET 命令傳遞。LOOKUP hint 和 OPTIONS hint 均不接受這類參數。
常用維表 JOIN 相關 TableConfig:
配置項 | 含義 | 預設值 |
| 非同步 lookup 的緩衝隊列容量(Job 級預設值) | 100 |
| 非同步 lookup 逾時時間(Job 級預設值) | 3min |
| 非同步 lookup 的輸出順序(Job 級預設值) |
|
-- 在作業開頭設定 Job 級預設值
SET 'table.exec.async-lookup.buffer-capacity' = '3072';
SET 'table.exec.async-lookup.timeout' = '180s';
INSERT INTO sink_table
SELECT ...
FROM src_table
LEFT JOIN dim_table FOR SYSTEM_TIME AS OF PROCTIME() AS d
ON src_table.key = d.key;當 SET 與 LOOKUP hint 同時設定同名參數時,單個 Join 操作上的 LOOKUP hint 優先順序更高。
通過 Join 策略 Hints 控制 Shuffle
維表 Join 策略 Hints 用於控制 Shuffle 行為,包括 SHUFFLE_HASH、REPLICATED_SHUFFLE_HASH 和 SKEW。維表 Cache 策略和聯結策略的適用情境如下。
Cache 策略 | SHUFFLE_HASH | REPLICATED_SHUFFLE_HASH(和 SKEW 等價) |
None | 不建議使用該聯結原則提示,主流會引入額外的網路開銷。 | 不建議使用該聯結原則提示,主流會引入額外的網路開銷。 |
LRU | 在維表尋找IO成為瓶頸時,建議考慮使用該聯結原則提示。當主流資料在Join Key上有時間局部性時,可以提高Cache命中率,減少IO請求數,從而提升總吞吐。 重要 主流會引入額外的網路開銷,當主流資料在Join Key上有傾斜,遇到效能瓶頸時,建議考慮REPLICATED_SHUFFLE_HASH。 | 在維表尋找IO成為瓶頸且主流資料在Join Key上有傾斜時,建議考慮該聯結原則提示。當主流資料在Join Key上有時間局部性時,可以提高Cache命中率,減少IO請求數,從而提升總吞吐。 |
ALL | 在維表記憶體使用量量成為瓶頸時,建議使用該聯結原則提示。記憶體使用量率可降低為1/並發度。 重要 主流會引入額外的網路開銷,當主流資料在Join Key上有傾斜,遇到效能瓶頸時,建議考慮REPLICATED_SHUFFLE_HASH。 | 在維表記憶體使用量量成為瓶頸且主流資料在Join Key上有傾斜時,建議使用該聯結原則提示。記憶體使用量率降低為分桶數/並發度。 |
當前 LOOKUP hint 的
shuffle選項已能覆蓋 SHUFFLE_HASH hint 功能,兩者同時使用時,會優先採納 LOOKUP hint 的shuffle選項。當前 LOOKUP hint 的
shuffle選項還未支援解決資料扭曲的功能,當和 REPLICATED_SHUFFLE_HASH、SKEW 同時使用時,會優先採納 REPLICATED_SHUFFLE_HASH、SKEW 對應的 shuffle 策略。
SHUFFLE_HASH
使用效果:在維表 Join 中使用 Shuffle Hash 策略,可將主流資料在 Join 之前根據 Join Key 做一次 shuffle。在使用 LRU Cache 策略時可提高 Cache 命中率、減少 IO 請求數;在使用 ALL Cache 策略時可減少記憶體使用量量。每個 SHUFFLE_HASH 聯結提示可指定多張維表。
使用限制:SHUFFLE_HASH 可減少記憶體開銷,但上遊資料需按 Join Key 做一次 shuffle,引入額外網路開銷,因此以下兩種情境不適合使用:
主流資料在 Join Key 上存在嚴重的資料扭曲。這種情境下使用 SHUFFLE_HASH,會因資料扭曲導致 Join 節點成為效能瓶頸,造成流作業嚴重反壓或批情境嚴重長尾,此時建議使用 REPLICATED_SHUFFLE_HASH。
維表資料較小,ALL Cache 策略載入無記憶體瓶頸時。這種情境下使用 SHUFFLE_HASH,節約的記憶體開銷和額外引入的網路開銷相比並不划算。
程式碼範例:
-- 只對維表 dim1 開啟 SHUFFLE_HASH 聯結
SELECT /*+ SHUFFLE_HASH(dim1) */ ...
-- 同時對維表 dim1、dim2 均開啟 SHUFFLE_HASH 聯結
SELECT /*+ SHUFFLE_HASH(dim1, dim2) */ ...
-- 維表 dim1 使用別名 D1 時,hint 中必須使用別名
SELECT /*+ SHUFFLE_HASH(D1) */ ...REPLICATED_SHUFFLE_HASH
使用效果:在維表 Join 中使用 Replicated Shuffle Hash 策略,效果基本與 SHUFFLE_HASH 一致,區別在於會將主流具有相同 key 的資料隨機打散到指定的 N 個並發上,可解決資料扭曲導致的效能瓶頸。每個 REPLICATED_SHUFFLE_HASH 聯結提示中可指定多張維表。
使用限制:
需配置傾斜資料分桶數量參數
table.exec.skew-join.replicate-num,預設值 16,取值不能大於維表聯結節點的並發。配置方法請參見《作業層級 SQL 調優》。當前不支援更新流,當主流是更新流時使用會報錯。
程式碼範例:
SELECT /*+ REPLICATED_SHUFFLE_HASH(dim1) */ ...SKEW
使用效果:當指定表存在資料扭曲時,最佳化器會在維表 Join 中使用 Replicated Shuffle Hash 策略(SKEW 是文法糖,底層用 Replicated Shuffle Hash 實現)。
使用限制:
每個 SKEW 提示只能指定 1 張表。
表名需為存在資料扭曲的主表名稱,而非維表名稱。
當前不支援更新流,當主流是更新流時使用會報錯。
程式碼範例:
SELECT /*+ SKEW(src) */ ...使用樣本
樣本一:基礎維表 JOIN
最基礎的寫法,使用 MySQL 維表補全 Kafka 流的欄位。無任何調優 hint。
CREATE TEMPORARY TABLE kafka_input (
id BIGINT,
name VARCHAR,
age BIGINT
) WITH (
'connector' = 'kafka',
'topic' = '<yourTopic>',
'properties.bootstrap.servers' = '<yourKafkaBrokers>',
'properties.group.id' = '<yourKafkaConsumerGroupId>',
'format' = 'csv'
);
CREATE TEMPORARY TABLE phoneNumber (
name VARCHAR,
phoneNumber BIGINT,
PRIMARY KEY (name) NOT ENFORCED
) WITH (
'connector' = 'mysql',
'hostname' = '<yourHostname>',
'port' = '3306',
'username' = '<yourUsername>',
'password' = '<yourPassword>',
'database-name' = '<yourDatabaseName>',
'table-name' = '<yourTableName>'
);
CREATE TEMPORARY TABLE result_infor (
id BIGINT,
phoneNumber BIGINT,
name VARCHAR
) WITH (
'connector' = 'blackhole'
);
INSERT INTO result_infor
SELECT
t.id,
w.phoneNumber,
t.name
FROM kafka_input AS t
JOIN phoneNumber FOR SYSTEM_TIME AS OF PROCTIME() AS w
ON t.name = w.name;樣本二:通過 OPTIONS hint 啟用 Cache 並通過 LOOKUP hint 配置非同步查詢
調優情境:維表 QPS 不足,需要在維表上開啟分區 Cache,同時啟用非同步 lookup 並增大緩衝隊列。
INSERT INTO user_behavior_wide
SELECT /*+ LOOKUP('table' = 't2', 'async' = 'true', 'capacity' = '3072') */
t1.member_id AS member_id,
t2.tag AS tag
FROM user_behavior_datagen AS t1
LEFT JOIN fluss.fluss.user_active_info /*+ OPTIONS(
'lookup.cache' = 'PARTIAL',
'lookup.partial-cache.max-rows' = '1000'
) */
FOR SYSTEM_TIME AS OF PROCTIME() AS t2
ON t1.member_id = t2.member_id;連接器選項(
lookup.cache、lookup.partial-cache.max-rows)通過OPTIONS()hint 寫在維表後,覆蓋建表時的 WITH 選項。Lookup 行為(
async、capacity)通過LOOKUP()hint 寫在 SELECT 之後。OPTIONS 與 LOOKUP 是兩類不同的 hint,必須分別書寫,混寫會導致 hint 靜默失效。
樣本三:通過 SET 調整 Job 級非同步 lookup 緩衝
調優情境:作業內多個維表 JOIN 都需要更大的非同步 lookup 緩衝,希望全域生效,並配合 LOOKUP hint 啟用 async。
SET 'table.exec.async-lookup.buffer-capacity' = '3072';
SET 'table.exec.async-lookup.timeout' = '180s';
INSERT INTO user_behavior_wide
SELECT /*+ LOOKUP('table' = 't2', 'async' = 'true') */
t1.member_id,
t2.tag
FROM user_behavior_datagen AS t1
LEFT JOIN fluss.fluss.user_active_info
FOR SYSTEM_TIME AS OF PROCTIME() AS t2
ON t1.member_id = t2.member_id;table.exec.async-lookup.buffer-capacity等參數為 Job 級 TableConfig,僅可通過 SET 設定,寫入 OPTIONS 或 LOOKUP hint 均不生效。若 LOOKUP hint 中未指定
capacity,則使用 SET 設定的全域值;若同時指定,則單個 JOIN 上的 hint 優先生效。