Delta Join(多流Join新范式)

更新时间:
复制 MD 格式

本文将为您介绍Delta Join的用法与实现。

Delta Join:基于Fluss的双流Join新方案

在实时数仓场景下,往往需要依赖多张实时数据表来构建统一的宽表。基于开源Flink + Kafka搭建的实时数仓,只能使用多个KafkaTopic通过Flink Join拼接形成一张大宽表,来实现无论哪个Topic发生更新时,总能对整个宽表完成更新。

由于Kafka本身并非面向分析场景设计,因此只能依托开源Flink的流式Join来实现这个场景,因为需要双边驱动更新,并在缓存全量上游数据,会导致开源Flink状态体积庞大,带来资源成本高、运维复杂、效率低等问题。

Fluss对比Kafka实现了Delta Join的能力。这是一种新型的 Join 方案,能够在保持双流 Join 语义的同时,将数据更新行为下沉到Fluss表中完成,从而显著降低Flink的资源消耗,提升作业的稳定性和执行效率。

Delta Join的优势

  • 无 Join State:省去冗余数据存储。

  • 低成本:仅依赖 Fluss 主键表和二级索引。

  • 更稳定高效 :避免大状态带来的性能瓶颈。

image

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 Keyuser_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。

image

相关参数

Flink 作业参数

参数名

默认值

说明

调整建议

table.optimizer.delta-join.strategy

NONE

是否改写为 Delta Join。
NONE:关闭。
AUTO:尝试改写,失败时自动回退为双流 Join。
EVENTUAL:强制改写,失败时报错。


说明

在实时计算引擎 VVR 11.5 前的版本中,EVENTUAL需要替换为FORCE

调试和生产都推荐 EVENTUAL,改写失败会直接报出原因,也能避免作业变更后静默退化成双流 Join。

table.optimizer.delta-join.ignore-non-unique-key-filter

false

是否跳过 Join 结果上非唯一键过滤条件的校验。

仅在业务能容忍结果表残留旧数据时设置为 true。该配置只跳过校验,不会改变过滤条件,也不会补发撤回消息。

table.exec.delta-join.cache-enabled

false

是否开启本地内存缓存,缓存命中时不再请求 Fluss。

内存压力不大时推荐设置为 true,可以明显减少对 Fluss 的点查请求。

table.exec.delta-join.left.cache-size

10000

缓存左表点查结果的 Key 数量。仅在开启缓存时生效。

右边每来一条数据,都会拿它的 Join Key 去查左表,查回的结果进入这块缓存,所以容量该设多大取决于右侧数据驱动点查的热点 Join Key 数量。

该参数会占用一部分内存,因此配置该参数时,也需要结合每条数据大小、TM内存大小来综合考虑。

在内存压力不大的情况下,可以先用默认值试一下。GC 频繁时适量减小。Fluss集群点查压力大时,适量增大。

table.exec.delta-join.right.cache-size

10000

缓存右表点查结果的 Key 数量。仅在开启缓存时生效。

左侧每来一条数据会去查右表。设置方法与左表缓存相同。

table.exec.async-lookup.buffer-capacity

100

每个 Delta 算子并发允许同时进行的异步点查请求数。

Fluss 集群压力和 TM CPU、内存压力都不大时调大,推荐调到千级别。

该参数按每个 Delta 算子的每个并发生效,作业在途请求总数约为并发数 × Delta 算子数量 × 该值,调大前请按总量评估 Fluss 侧的承载能力。出现内存紧张或 Fluss 侧过载时调小。

table.exec.async-lookup.timeout

3min

单次异步点查请求的超时时间。

点查偶发超时导致作业失败时调大,用于覆盖正常的尾部延迟。

Fluss 集群已经过载时不建议靠调大该值解决,超时请求会更长时间占用并发槽位,反而加重反压,此时应优先降低点查压力或扩容 Fluss。

Fluss 表参数

以下参数在 Fluss 建表的 WITH 子句中设置,也可以通过 SQL Hint 针对单个作业调整。

参数名

默认值

说明

调整建议

client.lookup.queue-size

25600

客户端等待处理的点查请求上限。

作业数据量大、点查排队明显时调大。TM 内存紧张时调小。

client.lookup.max-batch-size

128

合并为一个点查请求的最大条数。

请求量大、希望降低网络开销时调大。对延迟敏感时调小。

client.lookup.max-inflight-requests

128

当前同时处理的点查请求数上限。

提高点查并发时调大,需同步关注 Fluss 集群负载。

client.lookup.batch-timeout

100ms

等待批次攒满的最长时间,超时后立即发送。

对端到端延迟敏感时调小,希望提高批次填充率时调大。

通过 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;