投影增量构建
投影增量构建用于只处理上次成功更新后发生变化的数据,避免每次都重新计算完整视图。AIR 通过源表的增量字段识别变化范围,并将本次计算结果追加、合并或按受影响业务范围替换到已有投影数据中。
1. 核心概念
- 增量字段:源表中用于标识数据先后顺序的字段,例如自增 ID、批次号、创建时间或更新时间。
- 水位:上次成功构建后,投影结果中某个增量字段的最大值。
- 增量范围:本次构建时,满足
增量字段 > 上次水位的数据。 - 写入方式:可以通过
INSERT追加、通过MERGE按行更新或插入,也可以通过DELETE_INSERT替换本轮受影响业务范围的完整结果。
水位只在整次构建成功后推进
数据写入、最大值计算和构建结果保存都成功后,AIR 才会保存新水位。构建失败时,下次仍从最近一次成功构建的水位开始。
2. 适用范围
单表和多表视图均可以使用增量构建。对于多表视图,当前支持范围如下:
INSERT和MERGE写入方式支持由INNER JOIN、字段选择和过滤条件组成的视图。DELETE_INSERT写入方式额外支持LEFT JOIN;受影响键必须覆盖各增量源变化,并用于从完整视图重算对应业务范围。RIGHT JOIN和FULL JOIN暂不支持多表增量构建。- 可以为多张源表分别配置增量字段。
- 每张增量表在视图中只能出现一次。
- 每个增量字段都必须输出到视图最终结果,且输出字段名必须唯一。
- 增量字段不能经过会改变比较语义的复杂计算。
- 多表增量构建当前适用于使用 Spark 构建并写入 Iceberg 投影存储的场景。
支持的增量字段类型包括:
INTEGER、BIGINT等整数类型。FLOAT、DOUBLE、DECIMAL等数值类型。DATE、TIMESTAMP等日期时间类型。
增量字段的值应当随新数据单调增加。如果业务会把已有记录的增量字段改小,该记录可能无法被后续构建识别。
3. 开启前准备
3.1 配置源表增量字段
在视图依赖的来源表上配置增量字段。多表场景中,可以只为部分表配置,未配置的表在增量分支中按完整数据参与关联。
选择增量字段时,建议确保:
- 新增或变更的数据会产生更大的字段值。
- 字段不为空,或已明确空值数据不需要进入增量范围。
- 多张表的批次号不需要使用相同的数值,AIR 会分别维护水位。
3.2 在视图中输出增量字段
多表视图必须为各表的增量字段使用不同的最终输出名。例如:
SELECT
i.invoice_id,
i.batch_no AS invoice_batch_no,
p.payment_id,
p.batch_no AS payment_batch_no,
r.reconciliation_id,
r.batch_no AS reconciliation_batch_no,
c.customer_name
FROM invoice i
INNER JOIN payment p
ON i.payment_id = p.payment_id
INNER JOIN reconciliation r
ON p.reconciliation_id = r.reconciliation_id
INNER JOIN customer c
ON i.customer_id = c.customer_id;
上述示例中,发票表、支付表和对账表可以分别使用 batch_no 作为源增量字段,并在视图中输出为三个不同字段。客户表没有配置增量字段,但仍会参与每个分支的完整关联。
3.3 在投影中包含增量字段
创建明细投影时,必须在扫描字段中包含所有增量字段的最终输出列。否则 AIR 无法从投影结果计算各表的新水位。
4. 首次和后续构建
flowchart LR
A["创建并启用投影"] --> B["首次全量构建"]
B --> C["计算并保存各增量字段水位"]
C --> D["获取各表大于自身水位的数据"]
D --> E["计算本次增量结果"]
E --> F["INSERT、MERGE 或 DELETE_INSERT 写入"]
F --> G{"整次构建成功?"}
G -- "是" --> C
G -- "否" --> D
4.1 首次构建
一个新的增量构建周期开始时,AIR 会先执行一次全量构建:
- 按原始视图逻辑计算完整结果。
- 创建投影物理数据。
- 从构建成功的投影结果中,分别计算每个增量输出字段的最大值。
- 将这些最大值作为后续构建的初始水位。
首次构建不会向源表增加增量过滤条件。
4.2 后续构建
后续构建会读取最近一次成功构建的水位。单表场景会生成类似下面的条件:
多表场景会为每张增量表分别生成一个完整视图分支,详见下一节。
5. 多表增量的计算方式
假设发票表、支付表和对账表都配置了增量字段,上次水位分别为 I0、P0 和 R0。本次增量结果按以下逻辑计算:
-- 发票增量分支:增量发票关联完整支付和完整对账数据
<complete_view_query WHERE invoice.batch_no > I0>
UNION DISTINCT
-- 支付增量分支:增量支付关联完整发票和完整对账数据
<complete_view_query WHERE payment.batch_no > P0>
UNION DISTINCT
-- 对账增量分支:增量对账关联完整发票和完整支付数据
<complete_view_query WHERE reconciliation.batch_no > R0>
每个分支只过滤当前增量表,其他表使用当前完整数据。如果把三个增量条件同时放在一个分支中,只有三张表在同一轮都变化时才能命中,会遗漏单表迟到等情况。
5.1 关联数据迟到
假设发票 1099 已经存在,但它引用的支付记录尚未到达。由于使用 INNER JOIN,首次构建时不会得到该发票的投影结果。
支付记录后续到达时,支付增量分支会用它关联完整发票数据,从而生成发票 1099 的完整结果。
5.2 多个分支命中同一结果
如果一轮中同时新增了完整的发票、支付和对账关联链,同一结果可能在三个分支中都被计算到。AIR 使用 UNION DISTINCT 合并分支,对完全相同的结果进行去重。
UNION DISTINCT 只保证本次增量结果集内的去重
它不会自动判断业务主键,也不会清理历史投影中已有的重复数据。如果需要基于业务唯一键更新历史结果,应使用 MERGE。
5.3 A、B 增量表关联维度表 C
以下视图中,A、B 是增量表,C 是未配置增量字段的维度表:
SELECT
a.business_key,
a.updated_seq AS a_updated_seq,
b.updated_seq AS b_updated_seq,
c.dimension_name
FROM A a
INNER JOIN B b
ON a.b_id = b.id
LEFT JOIN C c
ON b.c_id = c.id;
该视图可以使用多表增量构建,但写入方式必须选择 DELETE_INSERT,并配置能够覆盖 A、B 两个增量源变化的受影响键。后续刷新按以下过程执行:
- 分别使用 A、B 的增量字段和最近一次成功水位识别变化数据。
- 从两个增量源的变化中计算并合并受影响键。
- 删除投影中这些受影响键对应的全部历史结果。
- 使用未添加增量过滤条件的完整 A、B、C 视图重新计算这些键的当前结果。
- 插入重算结果,并在整次构建成功后推进 A、B 各自的水位。
因此,当某条 B 数据的 c_id 从无匹配变为有匹配,或从有匹配变为无匹配时,只要该变化能够被 B 的增量字段识别,DELETE_INSERT 都会先撤回旧结果,再按当前 LEFT JOIN 结果完整重建受影响范围。
该场景不能使用 INSERT 或 MERGE
INSERT 只会追加本次结果,无法撤回受关联变化影响的旧结果。MERGE 只能处理本次 source 中仍然存在并能够命中目标的数据,无法安全覆盖关联失效、结果消失或一对多结果收缩。因此,包含 LEFT JOIN 的多表增量视图不会放开 INSERT 或 MERGE。
维度表 C 的变化不会自动触发刷新
C 未配置为增量表时,AIR 会在 A 或 B 触发刷新后读取 C 的当前完整数据,但 C 单独发生变化不会产生新的增量分支,也不会主动刷新投影。如果 C 的变化也必须及时反映到投影,应为 C 配置满足要求的增量字段,并确保受影响键覆盖 C 的变化,或者通过全量重建更新投影。
受影响键应选择稳定且能完整表示替换范围的业务键。Join Key 发生变化而增量数据又不包含变更前值时,必须通过“受影响键 SQL”或上游变更数据补充旧范围;否则 AIR 只能识别新范围,可能无法删除旧结果。
6. 选择 INSERT、MERGE 或 DELETE_INSERT
| 写入方式 | 适用场景 | 对已有数据的处理 | 主要限制 |
|---|---|---|---|
INSERT |
源数据只会新增,不会修改或撤回已有业务结果 | 直接追加本次增量结果 | 写入已成功但水位保存失败时,重试可能产生重复数据 |
MERGE |
已有结果的字段会更新,但变化后的记录仍会出现在本次增量结果中 | 按配置的 ON 条件匹配并更新、删除或插入单行 |
无法处理因逻辑删除、过滤失效、Join 失效等原因而不再出现在增量结果中的历史行 |
DELETE_INSERT |
需要撤回逻辑删除、过滤失效、Join 失效或一对多结果收缩产生的历史结果 | 先删除受影响键对应的全部旧结果,再插入该范围的完整当前结果 | 需要稳定且范围适当的替换键;当前适用于使用 Spark 构建并写入 Iceberg 的非分区投影 |
如果业务只会追加,优先使用 INSERT;如果已有行会更新但不会从结果中消失,使用 MERGE;只要业务变化可能使已有投影行不再满足视图条件,就应评估 DELETE_INSERT。完整选型案例请参见投影增量构建最佳实践。
6.1 INSERT 追加写入
在投影的高级配置中选择 INSERT。AIR 会产生类似下面的写入逻辑:
INSERT 不要求投影表存在主键,也不需要配置合并条件。
6.2 MERGE 合并写入
在投影的高级配置中选择 MERGE,并填写匹配和命中处理逻辑。例如,使用发票 ID 作为唯一条件:
当“错误数据处理”选择“继续插入”时,AIR 会为未匹配数据生成:
也可以使用多个字段组成联合唯一条件:
投影物理表不需要声明主键,但是用户配置的 ON 条件必须能在业务上唯一标识结果。如果同一次增量结果中存在多条数据匹配同一个目标行,MERGE 可能失败或得到非预期结果。
不会自动推导 MERGE 条件
AIR 不会根据源表主外键自动生成投影的 MERGE ON 条件。启用 MERGE 时必须提供有效的匹配表达式;配置缺失时构建会失败,不会自动改为 INSERT。
6.3 DELETE_INSERT 范围替换
DELETE_INSERT 用于处理“旧结果需要消失”的场景。例如,订单明细视图只保留 is_deleted = false 的记录。一条明细被逻辑删除后,它不会出现在本次最终增量结果中,因此普通 MERGE 没有源记录可以匹配并删除历史投影行。
使用 DELETE_INSERT 时,AIR 会执行以下操作:
- 根据最近一次成功水位识别本轮受影响的业务键,并对键集合去重。
- 删除投影中与这些键匹配的全部历史行。
- 使用完整视图的当前数据重新计算相同键范围。
- 插入该范围当前仍然存在的全部结果。
- 从更新后的投影结果计算并保存新水位。
假设订单 100 原来有两条明细,本轮其中一条被逻辑删除,并使用 order_id 作为替换键。AIR 会删除订单 100 的全部历史投影行,再插入订单 100 当前仍有效的所有明细。这样既能撤回已删除的明细,也不会丢失该订单下未发生变化的明细。
首次构建仍然是全量构建
新建投影或切换为 DELETE_INSERT 后,首次刷新使用完整视图创建投影数据,不执行范围删除。后续刷新才按受影响键执行删除和插入。
使用替换键
在“删除匹配模式”中选择“替换键”,再从投影最终输出字段中选择一个或多个字段。多个字段共同组成联合替换键。
替换键应满足以下条件:
- 能表示需要一起重算的最小业务一致性范围,例如订单明细场景中的
order_id。 - 存在于投影最终输出中,字段名称唯一,且可以从每个增量源通过明确的字段血缘推导。
- 值保持稳定。如果替换键或 Join Key 会被修改,但增量数据中没有旧值,AIR 只能识别新范围,无法删除旧范围。
“替换键”是推荐的默认模式。任一增量源无法完整推导替换键时,刷新会明确失败,不会忽略该增量源或缩小删除范围。此时可以改用“受影响键 SQL”。
使用受影响键 SQL
“受影响键 SQL”适用于替换键无法通过字段血缘自动推导的复杂传播关系。SQL 只负责描述“源数据变化会影响哪些投影业务键”,水位条件由 AIR 自动注入,不要在 SQL 中写死水位。
例如,订单、订单明细和商品表的变化都可能影响订单明细投影,可以返回受影响的 order_id:
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 必须满足以下要求:
- 只能包含一条只读查询,不能包含 DDL、DML 或多条语句。
- 不能引用当前投影的物理目标表。
- 至少返回一个字段;字段名必须唯一,并与投影输出字段同名且类型兼容。
- 必须覆盖配置中的每个增量源,使 AIR 能够为每个源分别应用最近一次成功水位。
- 全部输出字段共同组成联合替换键。SQL 范围过宽会增加删除和重算成本。
DELETE 和 INSERT 是顺序提交
当前 Iceberg 目标上的删除和插入是两个顺序提交。如果删除成功而插入失败,本次刷新不会保存新水位,后续重试仍会从最近一次成功水位重新计算并修复相同范围;在重试成功前,受影响范围可能暂时缺少数据。
7. 配置步骤
- 进入每张增量源表的详情页,配置增量字段。
- 创建或检查逻辑视图,确保所有增量字段都以唯一名称输出。
- 在视图详情页选择“投影”页签,新建明细投影。
- 在投影扫描字段中包含所有增量输出字段。
- 在“高级配置 → 增量配置”中选择
INSERT、MERGE或DELETE_INSERT。 - 根据写入方式补充配置:
MERGE:填写能够唯一匹配投影结果的ON条件及命中处理逻辑,并选择未匹配数据的处理方式。DELETE_INSERT:选择“替换键”或“受影响键 SQL”作为删除匹配模式,并完成对应配置。- 创建并启用投影,等待首次全量构建完成。
- 后续由调度或“更新投影”操作触发增量构建。
投影的基础配置和操作入口请参见投影功能。
8. 何时会执行全量构建
遇到以下任一情况时,AIR 会使用新的构建周期执行全量构建:
- 投影首次构建。
- 视图结构发生不兼容变化。
- 增量表集合、源增量字段、最终输出字段或字段类型发生变化。
- 写入方式在
INSERT、MERGE和DELETE_INSERT之间切换。 MERGE匹配表达式或未匹配数据处理方式发生变化。DELETE_INSERT的删除匹配模式、替换键或受影响键 SQL 发生变化。- 最近一次水位无法解析,或与当前增量配置不兼容。
如果多表视图不满足增量支持条件,本次刷新会整体按全量方式构建,不会只忽略其中一张增量表。常见原因包括:
INSERT或MERGE使用LEFT JOIN,或者任意写入方式使用RIGHT JOIN、FULL JOIN。- 使用聚合、窗口函数或集合运算。
- 同一张增量表在视图中被扫描多次,例如自关联。
- 增量字段未输出、输出名称重复、类型不支持,或增量条件无法下推到对应源表。
9. 使用限制与注意事项
- 物理删除必须有变化信号:直接从源表物理删除且没有 CDC 变更记录时,AIR 无法识别受影响键,三种写入方式都不会自动撤回历史结果。可使用包含旧键的业务变更表配合“受影响键 SQL”,或者执行全量刷新。
INSERT和MERGE不撤回消失结果:如果源数据变更后不再满足INNER JOIN或过滤条件,历史投影行不会被自动撤回;此类场景应使用DELETE_INSERT。多表视图需要使用LEFT JOIN时也必须选择DELETE_INSERT。- 替换键应保持稳定:替换键或 Join Key 发生变化但没有变更前值(before-image)时,
DELETE_INSERT无法识别旧范围。 - 替换范围影响构建成本:键的范围过大会增加完整视图重算和目标表改写成本;键过细则可能无法覆盖需要保持一致的关联结果。
- 没有跨源一致快照:多个增量分支读取的是构建时各源表的当前数据,不保证得到同一时刻的跨源快照。
- 分支数量会影响成本:每张增量表对应一个完整视图分支,增量表越多,源表扫描和
UNION DISTINCT去重开销越大。 - 水位来自投影结果:无法产生最终关联结果的较大批次值不会推进水位,因此可能在后续构建中被重复扫描。这可以保证关联数据迟到后仍能补齐结果。
INSERT不保证幂等:对重复写入敏感的业务,建议选择MERGE并使用稳定的业务唯一条件。
10. 验证构建结果
可在投影详情页的“更新记录”中检查构建结果:
- 确认首次构建为全量构建,后续构建为增量构建。
- 检查任务状态是否成功,以及输出数据量是否符合预期。
- 在 SQL 优化信息或 SQL 翻译详情中,检查每张增量表是否分别使用了
> 上次水位条件。 - 多表场景检查各分支是否使用
UNION合并,而不是将多个增量条件同时放入一个分支。 - 根据投影配置,确认最终写入使用
INSERT INTO、MERGE INTO,或按“受影响键准备 → DELETE → INSERT”的顺序执行。 - 对
MERGE场景,校验业务唯一键是否只保留一条结果。 - 对
DELETE_INSERT场景,检查受影响键数量、删除和插入阶段,并验证逻辑删除或失效关联对应的历史行已被撤回。
如果后续构建被执行为全量构建,应优先检查视图是否包含不支持的算子、增量字段是否仍在投影输出中,以及增量写入方式或匹配配置是否发生变化。