跳转至

投影增量构建最佳实践

本文通过一个订单明细宽表场景,逐步说明如何在 INSERTMERGEDELETE_INSERT 三种增量写入模式之间做选择,并给出配置、验证和运维建议。

1. 业务场景:订单明细宽表

某电商平台需要加速订单明细查询。逻辑视图关联以下三张表:

来源表 主要字段 变化方式
sales_order 订单 ID、客户 ID、订单状态、更新序号 新增订单,订单状态可能更新或取消
order_line 订单 ID、明细 ID、商品、数量、单价、逻辑删除标记、更新序号 新增或修改明细,明细可能被逻辑删除
product 商品编码、品类、更新序号 商品品类等属性可能更新

视图只保留未取消订单中的有效明细:

SELECT
    o.order_id,
    l.line_id,
    o.customer_id,
    o.status,
    l.sku,
    p.category_name,
    l.quantity,
    l.unit_price,
    l.quantity * l.unit_price AS line_amount,
    o.updated_seq AS order_updated_seq,
    l.updated_seq AS line_updated_seq,
    p.updated_seq AS product_updated_seq
FROM sales_order o
INNER JOIN order_line l
    ON o.order_id = l.order_id
INNER JOIN product p
    ON l.sku = p.sku
WHERE o.status <> 'CANCELLED'
  AND l.is_deleted = false;

在三张来源表中分别将 updated_seq 配置为增量字段,并确保对应字段以不同名称输出到视图和明细投影中。AIR 会为每张增量表独立维护水位。

按业务可能发生的最复杂变化选型

三种模式不是每批任务轮流使用的三个步骤。应根据同一个投影在整个生命周期中可能出现的数据变化选择一种模式并保持稳定。切换写入模式会开启新的构建周期,并先执行一次全量构建。

2. 第一步:数据只追加时使用 INSERT

2.1 业务假设

项目初期,订单进入来源表后不再修改:

  • 只会新增订单和订单明细。
  • 订单不会取消,明细不会修改或删除。
  • 商品属性不会变化。
  • 相同业务记录不会使用更大的 updated_seq 再次到达。

此时每轮增量结果都是历史结果之外的新行,不需要匹配已有投影数据。

2.2 推荐配置

  • 增量类型INSERT
  • 增量字段:选择三张表对应的增量输出字段
  • 业务唯一键:无需配置

AIR 将本轮增量结果直接追加到投影中。首次构建完成后,后续任务只处理大于各表最近成功水位的数据。

2.3 为什么 INSERT 最合适

INSERT 不需要对目标投影执行匹配,配置最简单,写入成本通常也最低。对于真正只追加的数据,没有必要为了防范不存在的更新使用更复杂的写入方式。

但“来源表有自增字段”不等于“业务只追加”。如果旧订单会以新的 updated_seq 再次写入,继续使用 INSERT 会保留新旧两份结果。

2.4 上线前验证

使用同一个 order_id 连续执行两轮测试:

  1. 首轮插入一条新订单和两条明细,刷新后应得到两条投影结果。
  2. 第二轮只插入另一张新订单,刷新后第一张订单仍为两条,且没有重复。
  3. 模拟一次任务重试,确认业务是否能够接受极端情况下重复追加的风险。

3. 第二步:已有行会更新时使用 MERGE

3.1 业务发生变化

随着业务发展,开始允许以下操作:

  • 修改明细数量或单价。
  • 修改商品品类。
  • 更新订单状态,但更新后的订单仍满足视图过滤条件。

这些变化会产生新的 updated_seq,而且变化后的记录仍会出现在本次增量结果中。例如订单 1001 的明细 10 数量从 1 改为 2,本轮增量结果仍包含 (1001, 10),只是字段值发生变化。

如果继续使用 INSERT,投影中会同时存在修改前和修改后的两行,因此应改用 MERGE

3.2 选择稳定的唯一匹配条件

该视图每一行表示一条订单明细,业务唯一键为 order_id + line_id。推荐配置:

  • 增量类型MERGE
  • MERGE 条件
ON t.order_id = s.order_id
AND t.line_id = s.line_id
WHEN MATCHED THEN UPDATE SET *
  • 错误数据处理:选择“作为新数据插入”,使新订单明细执行 WHEN NOT MATCHED THEN INSERT *

不要只使用 order_id 作为 MERGE ON 条件,因为一个订单包含多条明细,同一次增量结果中会有多条源记录匹配同一个目标范围,可能导致构建失败或错误更新。

3.3 MERGE 能解决什么

MERGE 适合“同一业务行仍然存在,只是字段值变化”的场景:

  • 新明细没有匹配项,执行插入。
  • 已有明细按联合唯一键命中,执行更新。
  • 任务使用旧水位重试时,可以再次匹配已写入的业务行,幂等性通常优于纯 INSERT

3.4 MERGE 的边界

MERGE 只能处理本次增量结果中存在的源记录。如果变化使一行不再满足视图逻辑,该行不会出现在 MERGE 的 source 中,也就无法命中历史投影行。

例如:

  • 明细 is_deletedfalse 变为 true
  • 订单状态变为 CANCELLED
  • 商品或关联键变化导致原有 INNER JOIN 不再成立。
  • 一对多关联结果从多行收缩为少行。

虽然可以在 WHEN MATCHED 中编写 DELETE,但前提仍然是 source 中存在一条能够匹配的记录。对于已经被视图过滤掉的结果,这个前提不成立。

3.5 上线前验证

  1. 修改订单 1001、明细 10 的数量并增加 line_updated_seq
  2. 刷新后确认 (1001, 10) 仍只有一行,且数量、金额均为新值。
  3. 同时新增一条明细,确认未匹配数据按配置插入。
  4. 检查同一轮增量结果中,order_id + line_id 是否始终唯一。

4. 第三步:结果可能消失时使用 DELETE_INSERT

4.1 业务继续演进

系统上线取消订单和逻辑删除明细功能。假设订单 1001 原来有两条投影结果:

order_id line_id is_deleted 是否在投影中
1001 10 false
1001 20 false

明细 20 被逻辑删除后,来源表中的 line_updated_seq 会增加,但这条记录不再满足 l.is_deleted = false。本次最终增量结果只剩明细 10,不会包含可以让 MERGE 删除明细 20 的源记录。

这时应使用 DELETE_INSERT,按业务范围重建结果。

4.2 推荐配置:按订单替换

  • 增量类型DELETE_INSERT
  • 删除匹配模式:替换键
  • 替换键order_id

每轮刷新时,AIR 会:

  1. 从各增量源变化中识别本轮受影响的 order_id
  2. 固定并去重受影响订单集合。
  3. 删除这些订单在投影中的全部历史明细。
  4. 使用完整视图重算这些订单。
  5. 插入每个订单当前仍有效的全部明细。

对于订单 1001,系统先删除明细 1020,再从完整视图插回仍有效的明细 10。最终既不会残留已删除的明细 20,也不会因为只处理变化记录而丢失未变化的明细 10

4.3 如何选择替换键粒度

替换键应表示业务上需要保持一致的最小重算范围。

候选替换键 评价
line_id 范围小,但无法自然覆盖订单级状态变化以及订单下需要一起保持一致的其他明细
order_id 能完整覆盖订单取消、明细删除和订单内一对多结果收缩,是本场景的推荐选择
customer_id 能覆盖客户下所有订单,但范围通常过大,会增加删除和重算成本

不要只追求最细的键,也不要为了简单而选择过粗的键。应以“源数据变化时,哪些结果必须一起撤回并按当前状态恢复”为判断标准。

4.4 何时使用受影响键 SQL

如果 AIR 无法从某个增量源自动推导 order_id,或者业务影响关系需要通过额外关联表达,可以把“删除匹配模式”改为“受影响键 SQL”。例如:

SELECT DISTINCT o.order_id
FROM sales_order o
INNER JOIN order_line l
    ON o.order_id = l.order_id
INNER JOIN product p
    ON l.sku = p.sku

编写时遵循以下原则:

  • SQL 只返回受影响的投影业务键,不返回完整投影结果。
  • 不要保留最终视图中的 status <> 'CANCELLED'is_deleted = false 过滤,否则刚被取消或删除的记录可能无法产生受影响键。
  • 不要写水位条件;AIR 会为每个配置的增量源分别注入水位。
  • 查询中的每个输出字段都必须与投影最终输出字段同名且类型兼容。
  • SQL 必须覆盖所有已选择的增量源,只能包含单条只读查询,不能引用投影物理目标表。

优先使用“替换键”。只有自动血缘无法准确表达影响范围时,才使用“受影响键 SQL”,并重点评估其返回范围和查询成本。

4.5 DELETE_INSERT 仍不能处理无信号的物理删除

DELETE_INSERT 需要先从增量数据中识别受影响键。如果一条来源记录被直接物理删除,且没有 CDC、变更日志或包含旧键的业务变更表,系统看不到大于水位的变化记录,因而无法撤回历史投影结果。

对于物理删除或替换键本身会变化的业务,应选择以下方案之一:

  • 保留逻辑删除记录并推进增量字段。
  • 使用包含变更前值(before-image)或旧业务键的变更表,并通过受影响键 SQL 映射替换范围。
  • 定期执行全量刷新进行兜底校准。

4.6 上线前验证

至少覆盖以下四类测试:

  1. 将一条明细标记为逻辑删除,确认历史投影行消失。
  2. 取消一个包含多条明细的订单,确认该订单的全部历史投影行消失。
  3. 只修改一个订单下的一条明细,确认该订单下未变化的兄弟明细仍然存在。
  4. 模拟删除成功、插入失败后重试,确认任务继续使用最近一次成功水位,并能修复完整订单范围。

5. 三种模式选型速查

flowchart TD
    A{"已有业务结果会变化吗?"}
    A -- "不会,只新增" --> B["使用 INSERT"]
    A -- "会" --> C{"变化后结果行仍会出现在<br/>本次增量结果中吗?"}
    C -- "始终会" --> D["使用 MERGE<br/>配置稳定唯一键"]
    C -- "可能不会" --> E["使用 DELETE_INSERT<br/>配置业务替换范围"]
判断问题 INSERT MERGE DELETE_INSERT
只新增数据 最佳选择 可以但通常没有必要 不建议,成本较高
更新已有结果 会产生重复 适合 可以,但会重算整个替换范围
逻辑删除或过滤失效 无法撤回 无法撤回 source 中已消失的行 适合
Join 失效或一对多收缩 无法撤回 通常无法完整撤回 适合
配置关键点 增量字段单调且业务只追加 ON 条件在业务上唯一稳定 替换键稳定、范围完整且不过宽
相对写入成本 较高,取决于受影响键数量和重算范围

6. 通用上线清单

6.1 配置检查

  • 每个新增或变化的业务记录都会产生更大的增量字段值。
  • 所有增量字段均以唯一名称输出到视图,并包含在明细投影中。
  • 多表视图使用 INSERTMERGE 时只包含 INNER JOIN;使用 LEFT JOIN 时选择 DELETE_INSERT 并配置覆盖完整的受影响键。
  • MERGE 的匹配条件能够唯一标识一条投影结果。
  • DELETE_INSERT 的替换键能够覆盖所有可能消失的历史结果,并且值保持稳定。
  • 选择 DELETE_INSERT 时,目标为使用 Spark 构建并写入 Iceberg 的非分区投影。

6.2 构建检查

  • 首次构建或写入模式变化后,先确认全量构建成功。
  • 后续刷新确认每个增量源都使用自己的最近成功水位。
  • 在投影详情页的“更新记录”中检查实际写入模式和各执行阶段。
  • DELETE_INSERT 重点关注受影响键数量、删除阶段和插入阶段;数量异常偏大时检查替换键粒度或受影响键 SQL。

6.3 结果检查

  • 对新增、更新、逻辑删除、过滤失效和 Join 失效分别准备测试数据。
  • 对比逻辑视图当前结果与刷新后的查询结果,重点检查重复行、残留行和缺失的兄弟结果。
  • 使用业务唯一键统计重复数据,不只比较总行数。
  • 在正式启用调度前至少验证一次失败重试路径。

更多产品行为和完整限制请参见投影增量构建