Delta Query 查询改写优化

更新时间:
复制 MD 格式

本文将为您介绍 Delta Query 的优化原理和相关参数。

什么是 Delta Query

Delta Query 是一类无状态化的流式查询改写。传统流算子为了增量维护结果,通常情况下需要在 Flink 状态中保留与历史数据规模成正比的中间数据,例如聚合累加器、Join 两侧的全量输入。随着作业运行时间增长,状态体积、资源占用和 Checkpoint 耗时都会同步上升,作业稳定性也会受到影响。

Delta Query 把这部分中间数据卸载到存储层,流算子几乎不保留状态。每收到一条数据,算子通过异步点查(Lookup)从源表取回所需明细,重新计算结果,并以幂等方式写入下游。源表使用支持毫秒级点查的高性能列式流存储 Fluss,Delta Query 作业可同时满足高吞吐与低延迟的需求。

核心优势

Delta Query 改写后的流作业,相比于有状态的流作业,有如下优势:

  • 将状态卸载至源表,Flink 作业省去冗余的状态存储,Checkpoint 体积和作业恢复时间相应下降,作业更加稳定。

  • 避免状态膨胀带来的性能瓶颈,节省 Flink 内存与 CPU。

  • 同一份源表明细可以被多个下游作业共享,消除冗余存储。

  • 排查异常结果时可以直接查询源表明细,便于验证结果、定位问题。

  • 表结构或查询发生变更后,可以基于同一份查询,先启动批作业处理历史数据,再通过指定位点的方式启动流作业,从而缩短断流时间,提升开发效率。

使用限制

  • 通过索引点查源表时,查询得到的明细数据行数需要可控。

Delta Query 每次收到变更,都会利用索引反查源表的明细数据。要查询的明细数据越多,单次点查返回的数据量和重算开销越大。

例如按大商家、超级用户这类高基数维度聚合全量历史明细时,单个分组可能达到数十万行,此时改写成 Delta Aggregate 后的开销会高于普通分组聚合,不建议改写。与此类似,在 Join 场景下,如果通过 Join Key 命中的单个索引从对侧查询出来的数据条数过多,也会对性能造成影响。

因此,启用 Delta Query 时需要评估单个索引的匹配规模。

  • 不同的 Delta Query 有不同的使用限制。只有满足了相应 Delta Query 的使用前提,才能成功启用 Delta Query,并且得到符合预期的计算结果。

具体请参考下文 Delta Join 使用限制Delta Aggregate 使用限制章节。

支持改写的查询类型

当前仅有部分查询支持改写优化为:

Delta 算子类型

面向的查询

最低引擎版本

Delta Join

流式双流 Join

VVR 8.0.11

Delta Aggregate

流式非窗口分组聚合

VVR 11.8

对源表和结果表的要求

对源表的要求

  1. 仅支持 Fluss 主键表作为源表

Delta Query 依赖源表提供毫秒级的点查能力,目前只支持 Fluss 主键表作为源表。在使用之前,需要先完成 Fluss 集群的购买和 Catalog 的创建,详情请参见开通流存储 Apache Fluss 版管理 Fluss Catalog

  1. 点查键必须完整命中源表的某个索引

改写后的算子按点查键反查源表明细。不同的查询依赖的点查键如下:

Delta 算子类型

点查键

Delta Join

Join 的等值条件,即 Join Key

Delta Aggregate

Aggregate 的分组字段,即 Group Key

点查键必须包含单个索引的全部字段,索引匹配不要求字段顺序一致。

Fluss 主键表当前支持两类索引:

  • 主键索引:建表声明主键(Primary Key)后即可使用,不需要额外配置。

  • 前缀索引:Fluss 主键表的分桶键(bucket.key)为主键(Primary Key)前缀,则该分桶键即为前缀索引键。如果该主键表是分区表,则分桶键联合分区键即为前缀索引键。

说明

例如 Fluss 主键表定义:PRIMARY KEY (user_id, order_id, order_data),且 'bucket.key' = 'user_id'

  • 如果该表为非分区表,则该表提供的索引包括主键索引user_id, order_id, order_data和前缀索引user_id

  • 如果该表为分区表,且order_data为分区键,则该表提供的索引包括主键索引user_id, order_id, order_data和前缀索引user_id, order_data

对结果表的要求

Delta Query 要求结果表能够按主键幂等地更新结果,结果表应当定义主键,并且连接器支持按主键更新。

结果表的主键还需要能够唯一标识一条查询结果:

  • Delta Join:主键需要唯一标识一条关联结果。两侧存在一对多或多对多关联时,主键通常需要同时包含两侧源表的主键字段,否则多条关联结果会互相覆盖。

  • Delta Aggregate:主键需要与 GROUP BY 字段一致,否则不同分组的结果会写入同一行。

说明

目前优化器并没有针对结果表主键的校验,请在建表时自行确认。

Delta Join

普通双流 Join 会在 Flink 状态里把两侧上游数据都存一份,并按 Join Key 建索引。Delta Join 不存这份数据,而是借助源表侧的索引,在任一侧变更到达后按 Join Key 反查另一侧的当前明细,输出关联结果。同一个 Join Key 上的多次变更按到达顺序串行处理,不同 Join Key 之间并行执行,因此不会出现较早的反查结果覆盖较晚结果的情况。

Delta Join 可以理解为双边驱动的维表 Join:左侧数据到达时查右表,右侧数据到达时查左表,再通过幂等更新写入结果表。

在业务可以容忍结果表无法物理删除数据、且确保源表主键和 Join Key 的关系不可变的前提下,Delta Join 的结果与普通双流 Join 一致。具体请参见业务逻辑限制

使用限制

引擎版本

  • 仅实时计算引擎 VVR 8.0.11 及以上版本支持 Delta Join。

  • VVR 8.0.11 及以上版本支持 INNER Join,VVR 11.3 及以上版本支持 LEFT/RIGHT/FULL OUTER Join,VVR 11.4 及以上版本支持级联 Delta Join。

查询结构

  • 仅支持改写流作业内的双流 Join 算子,不支持改写 Window Join、Interval Join、Temporal Join 和 Lookup Join。

  • Join 条件至少要有一组等值条件,且等值条件覆盖的列在两侧都要命中源表的单个索引。

  • Join Key 对应的索引列,从 Source 到 Join、以及级联 Join 之间不能参与任何计算。

  • 从 Source 到 Join、以及级联 Join 之间不能包含 Rand() 等非确定性函数,避免反查时无法复现相同结果。

  • 当前对 Delta Join 作业内的节点进行了限制,暂不支持白名单外的其他节点。

    • Source 到 Join 间、级联 Join 间,仅允许 Project 和 Filter 节点。源表上暂不支持定义 Watermark。定义 Watermark 会在 Source 之后引入 WatermarkAssigner 节点,不在当前白名单内。

    • Join 到 Sink 间仅允许 Project、Filter、Lookup(VVR 11.6 及以上版本支持)和 Union(VVR 11.8 及以上版本支持)节点。

    • 使用 BEGIN STATEMENT SET 同时写入多个 Sink 时:

      • VVR 11.8 以下版本,要求每一路写入都包含 Join 节点。

      • VVR 11.8 及以上版本,允许其中某一路不包含 Join 节点,但该路只能出现 Project、Filter、Lookup 和 Union 节点。

业务逻辑

  • 可以容忍结果表无法物理删除数据

Delta Join 只向下游发送 +I和 +U,不发送 -U-D。以过滤条件 WHERE l.amount + r.amount >= 100 为例,当某条 Join 记录从 120 掉到 80 时:

+U(..., 120) → 满足条件,写入结果表
+U(..., 80)  → 不满足条件,被过滤

由于没有 -U 撤旧结果,结果表里会残留 amount=120 的数据。

默认情况下,包含非唯一键过滤条件的 Query 会无法转化为 Delta Join。如果业务上能容忍数据残留时,可以设置 'table.optimizer.delta-join.ignore-non-unique-key-filter' = 'true' 跳过校验。

对于需要删除数据的场景,推荐使用逻辑删除,把过滤条件改写为标记列。以下面的查询为例,上游的每次变更都会覆盖写入结果表,is_delete 始终反映最新状态,下游查询时加上 WHERE is_delete = 0 即可。

INSERT INTO order_item_wide
SELECT
  o.merchant_id, o.order_id, o.item_id, i.item_name, o.amount,
  CASE WHEN o.amount >= 100 THEN 0 ELSE 1 END AS is_delete
FROM orders AS o
JOIN items AS i
  ON o.merchant_id = i.merchant_id AND o.item_id = i.item_id;

源表中的记录被删除时同理,删除消息不会传导到结果表,此前写入的关联结果会保留。需要让下游感知删除时,请使用上述的逻辑删除方案,由上游通过标记列表达删除状态。

  • 需要保证主键和 Join Key 的关系不可变,避免分布式乱序问题

例如某张源表,主键为user_id列,表上对应的 Join Key 为user_name,一旦user_id1user_nameJim的数据组成一条数据,不论后续user_id1对应的数据如何变化,user_name都不能更改为其他值,如Sam

使用示例

某电商平台需要把订单流和商品流关联成宽表,供下游查询。

步骤一:创建源表和结果表

两侧源表都要让 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;

Delta Aggregate

普通分组聚合会在 Flink 状态中为每个分组持续维护聚合累加器。Delta Aggregate 不长期维护各分组的聚合累加器,而是在收到数据变更后,用分组键异步点查源表,获取该分组的当前明细,基于这些明细重新计算聚合结果,聚合状态因此不再随分组数量增长。同一个分组键上的多次变更按到达顺序串行处理,不同分组键之间并行执行。

在业务可以容忍结果表无法物理删除数据、且确保主键和分组键的关系不可变的前提下,Delta Aggregate 的结果与普通分组聚合一致。具体请参见业务逻辑限制

使用限制

引擎版本

  • 仅实时计算引擎 VVR 11.8 及以上版本支持 Delta Aggregate。

查询结构

  • 查询结构

    • 仅支持改写流作业中的单层非窗口分组聚合,不支持窗口聚合和级联聚合。

    • Group Key 必须包含源表的单个索引。

    • Group Key 对应的索引列,从 Source 到 Aggregate 之间不能参与任何计算。

    • 从 Source 到 Aggregate 不能包含 Rand() 等非确定性函数,避免反查时无法复现相同结果。

    • 当前对 Delta Aggregate 作业内的节点进行了限制,暂不支持白名单外的其他节点。

      • Source 到 Aggregate 间,仅允许 Project、Filter 和 MiniBatchAssigner 节点。源表上暂不支持定义 Watermark。

      • Aggregate 到 Sink 间,仅允许 Project 和 Filter 节点。

      • 使用 BEGIN STATEMENT SET 同时写入多个 Sink 时每一路写入都必须包含 Aggregate 节点。

  • 聚合函数

    • 当前支持 SUMCOUNTAVGMINMAXLISTAGGFIRST_VALUELAST_VALUE

    • 暂不支持 DISTINCT、带 FILTER 子句的聚合、UDAF、Python 聚合函数,以及其他未列出的聚合函数。

说明

对于 LISTAGG、FIRST_VALUE、LAST_VALUE 等依赖输入顺序的聚合函数,Delta Aggregate 基于点查返回的明细重新计算,因此累加顺序与普通分组聚合按事件到达顺序累加不同。对结果顺序有确定性要求时,请谨慎使用这类函数。

业务逻辑

  • 可以容忍结果表无法物理删除数据

Delta Aggregate 只向下游发送+U-D,不发送-U。以过滤条件HAVING SUM(amount) >= 100 为例,当某个分组的SUM(amount)从 120 更新为 80 时:

+U(..., 120) → 满足条件,写入结果表
+U(..., 80)  → 不满足条件,被过滤

由于没有-U撤旧结果,结果表中仍会保留SUM(amount)=120的记录。

默认情况下,包含非唯一键过滤条件的 Query 会无法转化为 Delta Aggregate。如果业务上能容忍数据残留时,可以设置 'table.optimizer.delta-agg.ignore-non-unique-key-filter' = 'true' 跳过校验。

需要删除数据的场景,同样推荐使用逻辑删除,写法参见 Delta Join 的对应说明。

与 Delta Join 不同的是,当源表中的记录被删除时,Delta Aggregate 不丢弃删除消息,会产生正确的聚合结果且发送给下游。当分组内明细被删空时会向下游发送-D

  • 需要保证主键和 Group Key 的关系不可变,避免分布式乱序问题

例如某张源表,主键为user_id列,表上对应的 Group Key 为user_name,一旦user_id1user_nameJim的数据组成一条数据,不论后续user_id1对应的数据如何变化,user_name都不能更改为其他值,如Sam

使用示例

某电商平台需要按商家维度统计订单数和成交金额,结果写入汇总表,供商家看板直接读取。

步骤一:创建源表和结果表

本例的 GROUP BY merchant_id 没有覆盖源表主键 (merchant_id, order_id),但它是主键的前缀,因此声明为 Bucket Key,作为可用的前缀索引。

CREATE TABLE `my-catalog`.`my_db`.`orders` (
  merchant_id BIGINT,                    -- 商家 ID
  order_id BIGINT,                       -- 订单 ID
  amount DECIMAL(18, 2),                 -- 订单当前金额
  PRIMARY KEY (merchant_id, order_id) NOT ENFORCED
) WITH (
  'bucket.key' = 'merchant_id'           -- 按 merchant_id 反查该商家的全部订单
);

CREATE TABLE `my-catalog`.`my_db`.`merchant_summary` (
  merchant_id BIGINT,
  order_count BIGINT,                    -- 该商家当前订单数
  total_amount DECIMAL(38, 2),           -- 该商家当前成交总额,精度需高于源表字段以避免累加溢出
  PRIMARY KEY (merchant_id) NOT ENFORCED
);

步骤二:编写 Agg 作业

USE CATALOG `my-catalog`;
USE `my_db`;

-- 启用 Delta Aggregate 改写
SET 'table.optimizer.delta-agg.strategy' = 'EVENTUAL';

-- Delta Aggregate 语句本身不需要引入专用语法
INSERT INTO merchant_summary
SELECT
  merchant_id,
  COUNT(*) AS order_count,
  SUM(amount) AS total_amount
FROM orders
GROUP BY merchant_id;

步骤三:验证优化生效

作业部署上线、启动后,在“状态总览”界面的作业拓扑图中,如果看到如下的 Delta Aggregate 节点,则表明 Aggregate 已经成功被改写为 Delta Aggregate。

image

相关参数

Flink 作业参数

参数名

默认值

说明

调整建议

table.optimizer.delta-agg.strategy

NONE

是否改写为 Delta Aggregate。
NONE:关闭。
AUTO:尝试改写,失败时自动回退为普通聚合。
EVENTUAL:强制改写,失败时报错。


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

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

false

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

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

table.exec.delta-agg.cache-enabled

false

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

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

同一分组键被反复更新时收益明显。

table.exec.delta-agg.cache-size

10000

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

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

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

table.exec.mini-batch.enabled

false

是否开启微批处理。开启后,同一批次内相同分组键的变更会被合并,每个分组键在一个批次内最多发起一次点查。

源表点查次数偏高时开启。开启后端到端延迟会增加。

table.exec.mini-batch.allow-latency

0ms

微批处理最长等待时间。仅在开启微批处理时生效。

开启微批处理时必须设置为大于 0 的值。调大该值可以合并更多变更,减少点查源表的次数。

table.exec.mini-batch.size

-1

微批处理一个批次的最大条数。-1表示缓冲区满了之后自动触发微批处理。

仅在开启微批处理时生效。

维持默认值-1即可

table.exec.mini-batch.binary.memory

32MB

微批攒批的缓冲区大小。

仅在开启微批处理且微批大小为-1时生效。

维持默认值32MB即可

table.exec.async-lookup.buffer-capacity

100

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

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

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

table.exec.async-lookup.timeout

3min

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

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

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

Fluss 表参数

Delta Aggregate 的点查同样走 Fluss 客户端,可调整的参数与取值建议和 Delta Join 完全一致,请参见 Delta Join 章节的 Fluss 表参数。

常见问题

报错 Failed to perform delta-query optimization on the plan

这条报错表示强制改写失败。报错中会提示对应的策略参数,异常堆栈的 cause 里会说明具体原因,常见的有:

报错片段

原因

The SQL statement does not have any join node

查询里没有 Join,但 table.optimizer.delta-join.strategy 被设成了 EVENTUAL

The SQL statement does not have any group aggregate node

查询里没有分组聚合,但 table.optimizer.delta-agg.strategy 被设成了 EVENTUAL

The join key ... does not include all primary keys nor all fields from any index

Join Key 没有命中该侧源表的主键或任何索引。

The grouping key ... does not include all primary keys nor all fields from any index on the source table

GROUP BY 字段没有命中源表的任何索引。

The source table ... has neither a primary key nor an index

源表既没有主键也没有索引。

Unsupported node [...] encountered upstream/downstream

链路上出现了白名单外的节点,报错中会打印该节点及其输入。

注意,源表定义了 Watermark 时,会产生 WatermarkAssigner 节点,也可能命中这条报错。

Delta agg currently only supports the aggregate functions [...]

使用了不支持的聚合函数。

Table ... does not support async lookup

源表不支持异步点查。

Detected filters applied on non-unique keys

过滤条件引用了非唯一键字段。如果业务上可以容忍结果表无法物理删除数据,可以设置table.optimizer.delta-join.ignore-non-unique-key-filtertable.optimizer.delta-agg.ignore-non-unique-key-filter忽略该报错。

说明

不同实时计算引擎 VVR 版本中,报错信息可能会有变化,请以实际输出为准。