All Products
Search
Document Center

Realtime Compute for Apache Flink:Lookup join

Last Updated:Sep 11, 2026

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.enabled on 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.

  1. 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');
  1. Perform a lookup join with the insert-if-not-exists option. When a uid appears for the first time, Fluss automatically inserts it into the uid_mapping table and assigns an auto-incremented uid_int32 value.

-- 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_a appears for the first time, it is assigned uid_int32 = 1. This value is reused for subsequent appearances.

  • user_b and user_c are 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.