本文将为您介绍Delta Join的用法与实现。
Delta Join:基于Fluss的双流Join新方案
在实时数仓场景下,往往需要依赖多张实时数据表来构建统一的宽表。基于开源Flink + Kafka搭建的实时数仓,只能使用多个Kafka的Topic通过Flink Join拼接形成一张大宽表,来实现无论哪个Topic发生更新时,总能对整个宽表完成更新。
由于Kafka本身并非面向分析场景设计,因此只能依托开源Flink的流式Join来实现这个场景,因为需要双边驱动更新,并在缓存全量上游数据,会导致开源Flink状态体积庞大,带来资源成本高、运维复杂、效率低等问题。
Fluss对比Kafka实现了Delta Join的能力。这是一种新型的 Join 方案,能够在保持双流 Join 语义的同时,将数据更新行为下沉到Fluss表中完成,从而显著降低Flink的资源消耗,提升作业的稳定性和执行效率。
Delta Join的优势
无 Join State:省去冗余数据存储。
低成本:仅依赖 Fluss 主键表和二级索引。
更稳定高效 :避免大状态带来的性能瓶颈。
Delta Join使用限制
左右表必须为 Fluss 的主键表(支持分区表)。
Fluss 主键表的分桶键(Bucket Key)需要为主键(Primary Key)前缀。
Delta Join 的 Join Key 必须包含 Fluss 主键表定义的分桶键(Bucket Key),如果是分区表,则 Join Key 还需要包含分区键(Partition Key)。
说明例如定义:
PRIMARY KEY (user_id, order_id, order_data),且'bucket.key' = 'user_id'。Bucket key是主键的前缀,因此现在用户可以把
user_id当作二级索引进行高效的数据查找。对于查询:
JOIN users u ON o.user_id = u.user_id;此时Join Key为
user_id,包含了定义的分桶键,因此该查询会被优化成 Delta Join。
使用示例
某电商平台需要把订单流和商品流关联成宽表,供下游查询。
步骤一:创建源表和结果表
两侧源表都要让 Join Key 命中索引。本例的 Join Key 是 (merchant_id, item_id):在 orders 表中它是主键前缀,需要声明为 Bucket Key,作为可用的前缀索引;在 items 表中它就是主键,不需要额外配置。
CREATE TABLE `my-catalog`.`my_db`.`orders` (
merchant_id BIGINT, -- 商家 ID
item_id BIGINT, -- 商品 ID
order_id BIGINT, -- 订单 ID
amount DECIMAL(18, 2), -- 订单当前金额
PRIMARY KEY (merchant_id, item_id, order_id) NOT ENFORCED
) WITH (
'bucket.key' = 'merchant_id,item_id'
);
CREATE TABLE `my-catalog`.`my_db`.`items` (
merchant_id BIGINT, -- 商家 ID
item_id BIGINT, -- 商品 ID
item_name STRING, -- 商品名称
PRIMARY KEY (merchant_id, item_id) NOT ENFORCED -- Join Key 与主键一致
);
CREATE TABLE `my-catalog`.`my_db`.`order_item_wide` (
merchant_id BIGINT,
order_id BIGINT,
item_id BIGINT,
item_name STRING,
amount DECIMAL(18, 2),
PRIMARY KEY (merchant_id, order_id, item_id) NOT ENFORCED
);步骤二:编写 Join 作业
USE CATALOG `my-catalog`;
USE `my_db`;
-- 启用 Delta Join 改写
SET 'table.optimizer.delta-join.strategy' = 'EVENTUAL';
-- Delta Join 语句本身不需要引入专用语法
INSERT INTO order_item_wide
SELECT
o.merchant_id,
o.order_id,
o.item_id,
i.item_name,
o.amount
FROM orders AS o
JOIN items AS i
ON o.merchant_id = i.merchant_id AND o.item_id = i.item_id;步骤三:验证优化生效
作业部署上线、启动后,在“状态总览”界面的作业拓扑图中,如果看到如下的 Delta Join 节点,则表明双流 Join 已经成功被改写为 Delta Join。

相关参数
Flink 作业参数
参数名 | 默认值 | 说明 | 调整建议 |
|
| 是否改写为 Delta Join。 说明 在实时计算引擎 VVR 11.5 前的版本中, | 调试和生产都推荐 |
|
| 是否跳过 Join 结果上非唯一键过滤条件的校验。 | 仅在业务能容忍结果表残留旧数据时设置为 |
|
| 是否开启本地内存缓存,缓存命中时不再请求 Fluss。 | 内存压力不大时推荐设置为 |
|
| 缓存左表点查结果的 Key 数量。仅在开启缓存时生效。 | 右边每来一条数据,都会拿它的 Join Key 去查左表,查回的结果进入这块缓存,所以容量该设多大取决于右侧数据驱动点查的热点 Join Key 数量。 该参数会占用一部分内存,因此配置该参数时,也需要结合每条数据大小、TM内存大小来综合考虑。 在内存压力不大的情况下,可以先用默认值试一下。GC 频繁时适量减小。Fluss集群点查压力大时,适量增大。 |
|
| 缓存右表点查结果的 Key 数量。仅在开启缓存时生效。 | 左侧每来一条数据会去查右表。设置方法与左表缓存相同。 |
|
| 每个 Delta 算子并发允许同时进行的异步点查请求数。 | Fluss 集群压力和 TM CPU、内存压力都不大时调大,推荐调到千级别。 该参数按每个 Delta 算子的每个并发生效,作业在途请求总数约为并发数 × Delta 算子数量 × 该值,调大前请按总量评估 Fluss 侧的承载能力。出现内存紧张或 Fluss 侧过载时调小。 |
|
| 单次异步点查请求的超时时间。 | 点查偶发超时导致作业失败时调大,用于覆盖正常的尾部延迟。 Fluss 集群已经过载时不建议靠调大该值解决,超时请求会更长时间占用并发槽位,反而加重反压,此时应优先降低点查压力或扩容 Fluss。 |
Fluss 表参数
以下参数在 Fluss 建表的 WITH 子句中设置,也可以通过 SQL Hint 针对单个作业调整。
参数名 | 默认值 | 说明 | 调整建议 |
|
| 客户端等待处理的点查请求上限。 | 作业数据量大、点查排队明显时调大。TM 内存紧张时调小。 |
|
| 合并为一个点查请求的最大条数。 | 请求量大、希望降低网络开销时调大。对延迟敏感时调小。 |
|
| 当前同时处理的点查请求数上限。 | 提高点查并发时调大,需同步关注 Fluss 集群负载。 |
|
| 等待批次攒满的最长时间,超时后立即发送。 | 对端到端延迟敏感时调小,希望提高批次填充率时调大。 |
通过 Hint 调整时,Hint 需要写在被反查的表之后:
INSERT INTO order_item_wide
SELECT o.merchant_id, o.order_id, o.item_id, i.item_name, o.amount
FROM orders /*+ OPTIONS('client.lookup.queue-size' = '51200') */ AS o
JOIN items /*+ OPTIONS('client.lookup.queue-size' = '51200') */ AS i
ON o.merchant_id = i.merchant_id AND o.item_id = i.item_id;