This document explains how to perform a dimension table join.
Dimension table join (lookup join)
-
A lookup join is a specific implementation of a dimension table join. For more information, see Dimension Table JOIN Statement.
-
In Fluss, only primary key tables can be used as dimension tables. The join condition must include all columns of the primary key of the dimension table. If you set the bucket key as the join key when you create the table, the join condition can also be a prefix subset of the primary key. For more information, see Prefix lookup based on the bucket key.
-
If the dimension table is a partitioned primary key table, the join condition must also include the partition key.
-
By default, dimension table joins in Fluss use asynchronous mode for higher throughput. You can switch to synchronous mode by setting the
'lookup.async' = 'false'SQL hint.
SQL example for joining dimension tables
Create two dimension tables
USE CATALOG `fluss_catalog`;
USE `my_db`;
CREATE TABLE `orders` (
`o_orderkey` INT NOT NULL,
`o_custkey` INT NOT NULL,
`o_orderstatus` CHAR(1) NOT NULL,
`o_totalprice` DECIMAL(15, 2) NOT NULL,
`o_orderdate` DATE NOT NULL,
`o_orderpriority` CHAR(15) NOT NULL,
`o_clerk` CHAR(15) NOT NULL,
`o_shippriority` INT NOT NULL,
`o_comment` STRING NOT NULL,
`o_dt` STRING NOT NULL,
PRIMARY KEY (o_orderkey) NOT ENFORCED
);
CREATE TABLE `customer` (
`c_custkey` INT NOT NULL,
`c_name` STRING NOT NULL,
`c_address` STRING NOT NULL,
`c_nationkey` INT NOT NULL,
`c_phone` CHAR(15) NOT NULL,
`c_acctbal` DECIMAL(15, 2) NOT NULL,
`c_mktsegment` CHAR(10) NOT NULL,
`c_comment` STRING NOT NULL,
PRIMARY KEY (c_custkey) NOT ENFORCED
);
Join the two dimension tables
CREATE TEMPORARY TABLE `lookup_join_sink` (
`order_key` INT NOT NULL,
`order_totalprice` DECIMAL(15, 2) NOT NULL,
`customer_name` STRING NOT NULL,
`customer_address` STRING NOT NULL
) WITH ('connector' = 'blackhole');
-- Asynchronous mode (default)
INSERT INTO `lookup_join_sink`
SELECT `o`.`o_orderkey`, `o`.`o_totalprice`, `c`.`c_name`, `c`.`c_address`
FROM (SELECT `orders`.*, PROCTIME() AS `ptime` FROM `orders`) AS `o`
LEFT JOIN `customer` FOR SYSTEM_TIME AS OF `o`.`ptime` AS `c`
ON `o`.`o_custkey` = `c`.`c_custkey`;
Join in synchronous mode
-- Synchronous mode
INSERT INTO `lookup_join_sink`
SELECT `o`.`o_orderkey`, `o`.`o_totalprice`, `c`.`c_name`, `c`.`c_address`
FROM (SELECT `orders`.*, PROCTIME() AS `ptime` FROM `orders`) AS `o`
LEFT JOIN `customer` /*+ OPTIONS('lookup.async' = 'false') */ FOR SYSTEM_TIME AS OF `o`.`ptime` AS `c`
ON `o`.`o_custkey` = `c`.`c_custkey`;
Join a partitioned dimension table
If the dimension table is a partitioned primary key table, the join condition must include all primary key columns, that is, the primary key and the partition key.
-- Primary key: (c_custkey, dt). Partition key: dt
CREATE TABLE `customer_partitioned` (
`c_custkey` INT NOT NULL,
`c_name` STRING NOT NULL,
`c_address` STRING NOT NULL,
`c_nationkey` INT NOT NULL,
`c_phone` CHAR(15) NOT NULL,
`c_acctbal` DECIMAL(15, 2) NOT NULL,
`c_mktsegment` CHAR(10) NOT NULL,
`c_comment` STRING NOT NULL,
`dt` STRING NOT NULL,
PRIMARY KEY (`c_custkey`, `dt`) NOT ENFORCED
)
PARTITIONED BY (`dt`)
WITH (
'table.auto-partition.enabled' = 'true',
'table.auto-partition.time-unit' = 'year'
);
In addition to the primary key, the join condition must include the partition key.
INSERT INTO `lookup_join_sink`
SELECT `o`.`o_orderkey`, `o`.`o_totalprice`, `c`.`c_name`, `c`.`c_address`
FROM (SELECT `orders`.*, PROCTIME() AS `ptime` FROM `orders`) AS `o`
LEFT JOIN `customer_partitioned` FOR SYSTEM_TIME AS OF `o`.`ptime` AS `c`
ON `o`.`o_custkey` = `c`.`c_custkey` AND `o`.`o_dt` = `c`.`dt`;
For more information about partitioned tables, see Data partitioning.
Prefix lookup based on the bucket key
Overview
When a primary key table is used as a dimension table, if you set the bucket key as the join key when you create the table, the join condition can be a prefix subset of the primary key of the dimension table instead of the full primary key. This is a prefix lookup.
Prefix lookups use asynchronous mode by default for higher throughput. You can switch to synchronous mode by setting the 'lookup.async' = 'false' SQL hint.
Limits
-
Only Fluss primary key tables can be used as dimension tables for a prefix lookup.
-
You must set the bucket key as the join key when you create the table. For more information about how to set the bucket key, see Data bucketing.
-
The join condition must be a prefix subset of the primary key of the dimension table. For a partitioned primary key table, the join condition is a prefix of the primary key (excluding the partition key) plus the partition key.
-
A prefix lookup cannot be used together with
lookup.insert-if-not-exists.
Examples
Partitioned table
If the dimension table is a partitioned primary key table, the join condition is a prefix of the primary key (excluding the partition key) plus the partition key. In the following example, the primary key of the dimension table is (c_custkey, c_nationkey, dt), the bucket key is c_custkey, and the partition key is dt.
-- Primary key: (c_custkey, c_nationkey, dt). Bucket key: c_custkey. Partition key: dt
CREATE TABLE `customer_partitioned_with_bucket_key` (
`c_custkey` INT NOT NULL,
`c_name` STRING NOT NULL,
`c_address` STRING NOT NULL,
`c_nationkey` INT NOT NULL,
`c_phone` CHAR(15) NOT NULL,
`c_acctbal` DECIMAL(15, 2) NOT NULL,
`c_mktsegment` CHAR(10) NOT NULL,
`c_comment` STRING NOT NULL,
`dt` STRING NOT NULL,
PRIMARY KEY (`c_custkey`, `c_nationkey`, `dt`) NOT ENFORCED
)
PARTITIONED BY (`dt`)
WITH (
'bucket.key' = 'c_custkey',
'table.auto-partition.enabled' = 'true',
'table.auto-partition.time-unit' = 'year'
);
-- The join condition is a prefix of the primary key (excluding the partition key) plus the partition key
INSERT INTO `lookup_join_sink`
SELECT `o`.`o_orderkey`, `o`.`o_totalprice`, `c`.`c_name`, `c`.`c_address`
FROM (SELECT `orders`.*, PROCTIME() AS `ptime` FROM `orders`) AS `o`
LEFT JOIN `customer_partitioned_with_bucket_key` FOR SYSTEM_TIME AS OF `o`.`ptime` AS `c`
ON `o`.`o_custkey` = `c`.`c_custkey` AND `o`.`o_dt` = `c`.`dt`;
Non-partitioned table
In the following example, the primary key of the dimension table is (c_custkey, c_nationkey), the bucket key is c_custkey, and the join condition uses only the primary key prefix c_custkey.
-- Primary key: (c_custkey, c_nationkey). Bucket key: c_custkey
CREATE TABLE `customer_with_bucket_key` (
`c_custkey` INT NOT NULL,
`c_name` STRING NOT NULL,
`c_address` STRING NOT NULL,
`c_nationkey` INT NOT NULL,
`c_phone` CHAR(15) NOT NULL,
`c_acctbal` DECIMAL(15, 2) NOT NULL,
`c_mktsegment` CHAR(10) NOT NULL,
`c_comment` STRING NOT NULL,
PRIMARY KEY (`c_custkey`, `c_nationkey`) NOT ENFORCED
) WITH (
'bucket.key' = 'c_custkey'
);
-- Prefix lookup in asynchronous mode (default)
INSERT INTO `lookup_join_sink`
SELECT `o`.`o_orderkey`, `o`.`o_totalprice`, `c`.`c_name`, `c`.`c_address`
FROM (SELECT `orders`.*, PROCTIME() AS `ptime` FROM `orders`) AS `o`
LEFT JOIN `customer_with_bucket_key` FOR SYSTEM_TIME AS OF `o`.`ptime` AS `c`
ON `o`.`o_custkey` = `c`.`c_custkey`;
-- Prefix lookup in synchronous mode
INSERT INTO `lookup_join_sink`
SELECT `o`.`o_orderkey`, `o`.`o_totalprice`, `c`.`c_name`, `c`.`c_address`
FROM (SELECT `orders`.*, PROCTIME() AS `ptime` FROM `orders`) AS `o`
LEFT JOIN `customer_with_bucket_key` /*+ OPTIONS('lookup.async' = 'false') */ FOR SYSTEM_TIME AS OF `o`.`ptime` AS `c`
ON `o`.`o_custkey` = `c`.`c_custkey`;
Historical partition lookup
Overview
Auto partitioning removes expired partitions based on the retention policy. After a partition is removed, a lookup join that references the partition can no longer find its rows, even if the data has been tiered to Paimon.
After historical partition lookup is enabled, primary key lookups fall back to Paimon when the original partition no longer exists. To enable this feature, set table.datalake.historical-partition.enabled on the dimension table.
Limits
-
This feature is disabled by default. Enable it by setting
table.datalake.historical-partition.enabledon the dimension table. -
Only Paimon primary key tables with auto partitioning enabled and exactly one partition key are supported.
-
After the feature is enabled, the Coordinator creates and retains the
__historical__system partition to route lookups to Paimon. Disabling the option removes the system partition. -
A lookup client uses the table configuration captured when the lookuper is created to decide whether to fall back. After you change the configuration, restart the lookup jobs that need to query historical partitions.
Examples
Enable historical partition lookup on the dimension table.
ALTER TABLE `customer_partitioned_with_bucket_key` SET (
'table.datalake.historical-partition.enabled' = 'true'
);
For more information about the parameter, see Connector parameters.
Insert if not exists
Overview
During a lookup join, if a lookup key is not found in the dimension table, the join is skipped by default. For a LEFT JOIN, columns from the dimension table are NULL. When the lookup.insert-if-not-exists option is enabled and a lookup miss occurs, Fluss automatically inserts a new row. This row uses the lookup key as its primary key and is returned as the join result.
This feature is particularly useful for automatically building a dictionary table (which must be a primary key table) during stream processing. A typical use case is to map high-cardinality string identifiers, such as user IDs or device IDs, to compact integer IDs for efficient processing by the downstream Aggregation Merge Engine.
Limits
-
Only primary key lookups are supported; prefix lookups are not.
-
The dimension table cannot contain non-nullable columns, other than primary key and auto-increment columns, because Fluss cannot populate their values during automatic insertion.
-
Enable this feature with the following SQL hint:
/*+ OPTIONS('lookup.insert-if-not-exists' = 'true') */
Examples
The following example demonstrates how to automatically build a UID dictionary table in a lookup join.
-
Create a dictionary table with an auto-increment column.
CREATE TABLE uid_mapping (
uid VARCHAR NOT NULL,
uid_int32 INT,PRIMARY KEY (uid) NOT ENFORCED
) WITH ('auto-increment.fields' = 'uid_int32','bucket.num' = '1');
-
Perform a lookup join with the
insert-if-not-existsoption. When auidappears for the first time, Fluss automatically inserts it into theuid_mappingtable and assigns an auto-incrementeduid_int32value.
-- Registers UIDs from the ods_events stream table in the uid_mapping dictionary table
-- and retrieves the corresponding integer ID (uid_int32).
SELECT
ods.country,
ods.prov,
ods.city,
ods.ymd,
ods.uid,
dim.uid_int32
FROM ods_events AS ods
JOIN uid_mapping /*+ OPTIONS('lookup.insert-if-not-exists' = 'true') */FOR SYSTEM_TIME AS OF ods.proctime AS dim
ON dim.uid = ods.uid;
The ods_events table contains the following data:
|
Country |
Prov |
City |
Ymd |
Uid |
|
CN |
Beijing |
Haidian |
2025-01-01 |
user_a |
|
CN |
Shanghai |
Pudong |
2025-01-02 |
user_b |
|
US |
California |
LA |
2025-01-03 |
user_a |
|
JP |
Tokyo |
Shibuya |
2025-01-04 |
user_c |
The join result is as follows:
|
Country |
Prov |
City |
Ymd |
Uid |
Uid_int32 |
|
CN |
Beijing |
Haidian |
2025-01-01 |
user_a |
1 |
|
CN |
Shanghai |
Pudong |
2025-01-02 |
user_b |
2 |
|
US |
California |
LA |
2025-01-03 |
user_a |
1 |
|
JP |
Tokyo |
Shibuya |
2025-01-04 |
user_c |
3 |
-
When
user_aappears for the first time, it is assigneduid_int32 = 1. This value is reused for subsequent appearances. -
user_banduser_care each assigned a new auto-incremented ID.
After the job runs, the uid_mapping dictionary table contains the following data:
|
Uid |
Uid_int32 |
|
user_a |
1 |
|
user_b |
2 |
|
user_c |
3 |
References
-
For the configuration options supported by dimension table joins, see Connector parameters.
-
For more information about how to set the bucket key, see Data bucketing.
-
For more information about auto partitioning and partition management strategies, see Data partitioning.