全部產品
Search
文件中心

Realtime Compute for Apache Flink:維表JOIN語句

更新時間:May 29, 2026

對於每條流式資料,可以關聯一個外部維表資料來源,為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.cachelookup.partial-cache.max-rowslookup.partial-cache.expire-after-write

建表語句的 WITH (...);或在 SQL 中用 /*+ OPTIONS('key'='value') */ 覆蓋

單張表

LOOKUP hint 選項

tableasynccapacitytimeoutoutput-moderetry-predicateretry-strategyfixed-delaymax-attemptsshuffle

/*+ LOOKUP('table'='dim', 'async'='true', 'capacity'='100') */

單個 Join 操作

Flink Job 級 TableConfig

table.exec.*table.optimizer.* 開頭的全域配置

作業開頭使用 SET 'table.exec.async-lookup.buffer-capacity' = '3072';

整個作業

重要
  • 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.cachelookup.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 中必須使用別名。

支援的選項

選項

含義

取值

table

指定 hint 作用的維表名或別名

表名/別名字串

async

是否啟用非同步 lookup

true / false

output-mode

非同步 lookup 的輸出順序

ordered / allow_unordered

capacity

非同步 lookup 緩衝隊列容量

整數

timeout

非同步 lookup 逾時時間

時間長度(如 180s

retry-predicate

觸發重試的條件

當前僅支援 lookup_miss

retry-strategy

重試策略

當前僅支援 fixed_delay

fixed-delay

固定稍候再試

時間長度(如 10s

max-attempts

最大重試次數

整數

shuffle

是否在維表 Join 前做 shuffle

true / false

shuffle 選項行為

shuffle 選項影響維表 Join 的 shuffle 策略,不同情境表現如下。

情境

聯結策略

不配置 'shuffle' = 'true' 選項

使用引擎預設的 shuffle 策略。

不配置 'shuffle' = 'true' 選項,且維表連接器不提供自訂聯結策略

使用引擎預設的 shuffle 策略。

配置 'shuffle' = 'true' 選項,且維表連接器不提供自訂聯結策略

預設使用 SHUFFLE_HASH 策略,含義請參見《SHUFFLE_HASH》。

配置 'shuffle' = 'true' 選項,且維表連接器提供自訂聯結策略

使用表連接器的自訂 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:

配置項

含義

預設值

table.exec.async-lookup.buffer-capacity

非同步 lookup 的緩衝隊列容量(Job 級預設值)

100

table.exec.async-lookup.timeout

非同步 lookup 逾時時間(Job 級預設值)

3min

table.exec.async-lookup.output-mode

非同步 lookup 的輸出順序(Job 級預設值)

ordered

-- 在作業開頭設定 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.cachelookup.partial-cache.max-rows)通過 OPTIONS() hint 寫在維表後,覆蓋建表時的 WITH 選項。

  • Lookup 行為asynccapacity)通過 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 優先生效。