跳转至

投影增量构建

投影增量构建用于只处理上次成功更新后发生变化的数据,避免每次都重新计算完整视图。AIR 通过源表的增量字段识别变化范围,并将本次计算结果追加、合并或按受影响业务范围替换到已有投影数据中。

1. 核心概念

  • 增量字段:源表中用于标识数据先后顺序的字段,例如自增 ID、批次号、创建时间或更新时间。
  • 水位:上次成功构建后,投影结果中某个增量字段的最大值。
  • 增量范围:本次构建时,满足 增量字段 > 上次水位 的数据。
  • 写入方式:可以通过 INSERT 追加、通过 MERGE 按行更新或插入,也可以通过 DELETE_INSERT 替换本轮受影响业务范围的完整结果。

水位只在整次构建成功后推进

数据写入、最大值计算和构建结果保存都成功后,AIR 才会保存新水位。构建失败时,下次仍从最近一次成功构建的水位开始。

2. 适用范围

单表和多表视图均可以使用增量构建。对于多表视图,当前支持范围如下:

  • INSERTMERGE 写入方式支持由 INNER JOIN、字段选择和过滤条件组成的视图。
  • DELETE_INSERT 写入方式额外支持 LEFT JOIN;受影响键必须覆盖各增量源变化,并用于从完整视图重算对应业务范围。
  • RIGHT JOINFULL JOIN 暂不支持多表增量构建。
  • 可以为多张源表分别配置增量字段。
  • 每张增量表在视图中只能出现一次。
  • 每个增量字段都必须输出到视图最终结果,且输出字段名必须唯一。
  • 增量字段不能经过会改变比较语义的复杂计算。
  • 多表增量构建当前适用于使用 Spark 构建并写入 Iceberg 投影存储的场景。

支持的增量字段类型包括:

  • INTEGERBIGINT等整数类型。
  • FLOATDOUBLEDECIMAL 等数值类型。
  • DATETIMESTAMP 等日期时间类型。

增量字段的值应当随新数据单调增加。如果业务会把已有记录的增量字段改小,该记录可能无法被后续构建识别。

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 会先执行一次全量构建:

  1. 按原始视图逻辑计算完整结果。
  2. 创建投影物理数据。
  3. 从构建成功的投影结果中,分别计算每个增量输出字段的最大值。
  4. 将这些最大值作为后续构建的初始水位。

首次构建不会向源表增加增量过滤条件。

4.2 后续构建

后续构建会读取最近一次成功构建的水位。单表场景会生成类似下面的条件:

WHERE batch_no > <last_watermark>

多表场景会为每张增量表分别生成一个完整视图分支,详见下一节。

5. 多表增量的计算方式

假设发票表、支付表和对账表都配置了增量字段,上次水位分别为 I0P0R0。本次增量结果按以下逻辑计算:

-- 发票增量分支:增量发票关联完整支付和完整对账数据
<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 两个增量源变化的受影响键。后续刷新按以下过程执行:

  1. 分别使用 A、B 的增量字段和最近一次成功水位识别变化数据。
  2. 从两个增量源的变化中计算并合并受影响键。
  3. 删除投影中这些受影响键对应的全部历史结果。
  4. 使用未添加增量过滤条件的完整 A、B、C 视图重新计算这些键的当前结果。
  5. 插入重算结果,并在整次构建成功后推进 A、B 各自的水位。

因此,当某条 B 数据的 c_id 从无匹配变为有匹配,或从有匹配变为无匹配时,只要该变化能够被 B 的增量字段识别,DELETE_INSERT 都会先撤回旧结果,再按当前 LEFT JOIN 结果完整重建受影响范围。

该场景不能使用 INSERT 或 MERGE

INSERT 只会追加本次结果,无法撤回受关联变化影响的旧结果。MERGE 只能处理本次 source 中仍然存在并能够命中目标的数据,无法安全覆盖关联失效、结果消失或一对多结果收缩。因此,包含 LEFT JOIN 的多表增量视图不会放开 INSERTMERGE

维度表 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 INTO <projection_table>
<incremental_result>;

INSERT 不要求投影表存在主键,也不需要配置合并条件。

6.2 MERGE 合并写入

在投影的高级配置中选择 MERGE,并填写匹配和命中处理逻辑。例如,使用发票 ID 作为唯一条件:

ON t.invoice_id = s.invoice_id
WHEN MATCHED THEN UPDATE SET *

当“错误数据处理”选择“继续插入”时,AIR 会为未匹配数据生成:

WHEN NOT MATCHED THEN INSERT *

也可以使用多个字段组成联合唯一条件:

ON t.invoice_id = s.invoice_id
AND t.customer_id = s.customer_id
WHEN MATCHED THEN UPDATE SET *

投影物理表不需要声明主键,但是用户配置的 ON 条件必须能在业务上唯一标识结果。如果同一次增量结果中存在多条数据匹配同一个目标行,MERGE 可能失败或得到非预期结果。

不会自动推导 MERGE 条件

AIR 不会根据源表主外键自动生成投影的 MERGE ON 条件。启用 MERGE 时必须提供有效的匹配表达式;配置缺失时构建会失败,不会自动改为 INSERT

6.3 DELETE_INSERT 范围替换

DELETE_INSERT 用于处理“旧结果需要消失”的场景。例如,订单明细视图只保留 is_deleted = false 的记录。一条明细被逻辑删除后,它不会出现在本次最终增量结果中,因此普通 MERGE 没有源记录可以匹配并删除历史投影行。

使用 DELETE_INSERT 时,AIR 会执行以下操作:

  1. 根据最近一次成功水位识别本轮受影响的业务键,并对键集合去重。
  2. 删除投影中与这些键匹配的全部历史行。
  3. 使用完整视图的当前数据重新计算相同键范围。
  4. 插入该范围当前仍然存在的全部结果。
  5. 从更新后的投影结果计算并保存新水位。

假设订单 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. 配置步骤

  1. 进入每张增量源表的详情页,配置增量字段。
  2. 创建或检查逻辑视图,确保所有增量字段都以唯一名称输出。
  3. 在视图详情页选择“投影”页签,新建明细投影。
  4. 在投影扫描字段中包含所有增量输出字段。
  5. 在“高级配置 → 增量配置”中选择 INSERTMERGEDELETE_INSERT
  6. 根据写入方式补充配置:
  7. MERGE:填写能够唯一匹配投影结果的 ON 条件及命中处理逻辑,并选择未匹配数据的处理方式。
  8. DELETE_INSERT:选择“替换键”或“受影响键 SQL”作为删除匹配模式,并完成对应配置。
  9. 创建并启用投影,等待首次全量构建完成。
  10. 后续由调度或“更新投影”操作触发增量构建。

投影的基础配置和操作入口请参见投影功能

8. 何时会执行全量构建

遇到以下任一情况时,AIR 会使用新的构建周期执行全量构建:

  • 投影首次构建。
  • 视图结构发生不兼容变化。
  • 增量表集合、源增量字段、最终输出字段或字段类型发生变化。
  • 写入方式在 INSERTMERGEDELETE_INSERT 之间切换。
  • MERGE 匹配表达式或未匹配数据处理方式发生变化。
  • DELETE_INSERT 的删除匹配模式、替换键或受影响键 SQL 发生变化。
  • 最近一次水位无法解析,或与当前增量配置不兼容。

如果多表视图不满足增量支持条件,本次刷新会整体按全量方式构建,不会只忽略其中一张增量表。常见原因包括:

  • INSERTMERGE 使用 LEFT JOIN,或者任意写入方式使用 RIGHT JOINFULL JOIN
  • 使用聚合、窗口函数或集合运算。
  • 同一张增量表在视图中被扫描多次,例如自关联。
  • 增量字段未输出、输出名称重复、类型不支持,或增量条件无法下推到对应源表。

9. 使用限制与注意事项

  • 物理删除必须有变化信号:直接从源表物理删除且没有 CDC 变更记录时,AIR 无法识别受影响键,三种写入方式都不会自动撤回历史结果。可使用包含旧键的业务变更表配合“受影响键 SQL”,或者执行全量刷新。
  • INSERTMERGE 不撤回消失结果:如果源数据变更后不再满足 INNER JOIN 或过滤条件,历史投影行不会被自动撤回;此类场景应使用 DELETE_INSERT。多表视图需要使用 LEFT JOIN 时也必须选择 DELETE_INSERT
  • 替换键应保持稳定:替换键或 Join Key 发生变化但没有变更前值(before-image)时,DELETE_INSERT 无法识别旧范围。
  • 替换范围影响构建成本:键的范围过大会增加完整视图重算和目标表改写成本;键过细则可能无法覆盖需要保持一致的关联结果。
  • 没有跨源一致快照:多个增量分支读取的是构建时各源表的当前数据,不保证得到同一时刻的跨源快照。
  • 分支数量会影响成本:每张增量表对应一个完整视图分支,增量表越多,源表扫描和 UNION DISTINCT 去重开销越大。
  • 水位来自投影结果:无法产生最终关联结果的较大批次值不会推进水位,因此可能在后续构建中被重复扫描。这可以保证关联数据迟到后仍能补齐结果。
  • INSERT 不保证幂等:对重复写入敏感的业务,建议选择 MERGE 并使用稳定的业务唯一条件。

10. 验证构建结果

可在投影详情页的“更新记录”中检查构建结果:

  1. 确认首次构建为全量构建,后续构建为增量构建。
  2. 检查任务状态是否成功,以及输出数据量是否符合预期。
  3. 在 SQL 优化信息或 SQL 翻译详情中,检查每张增量表是否分别使用了 > 上次水位 条件。
  4. 多表场景检查各分支是否使用 UNION 合并,而不是将多个增量条件同时放入一个分支。
  5. 根据投影配置,确认最终写入使用 INSERT INTOMERGE INTO,或按“受影响键准备 → DELETE → INSERT”的顺序执行。
  6. MERGE 场景,校验业务唯一键是否只保留一条结果。
  7. DELETE_INSERT 场景,检查受影响键数量、删除和插入阶段,并验证逻辑删除或失效关联对应的历史行已被撤回。

如果后续构建被执行为全量构建,应优先检查视图是否包含不支持的算子、增量字段是否仍在投影输出中,以及增量写入方式或匹配配置是否发生变化。