跳转至

投影增量构建

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

1. 核心概念

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

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

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

2. 适用范围

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

  • 视图使用 INNER 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 写入"]
    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

6. 选择 INSERT 或 MERGE

写入方式 适用场景 对已有数据的处理 重试特性
INSERT 源数据只会新增,不会修改已有业务记录 直接追加本次增量结果 写入已成功但水位保存失败时,重试可能产生重复数据
MERGE 需要按业务唯一键更新已有数据,或需要增强重试幂等性 按外部配置的 ON 条件匹配并执行更新、删除或插入 ON 条件在本次数据中唯一时,重试可匹配已写入的结果

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

7. 配置步骤

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

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

8. 何时会执行全量构建

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

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

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

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

9. 使用限制与注意事项

  • 不处理物理删除:源表中已删除的行不会因此自动从历史投影中删除。
  • 不撤回失效关联结果:如果源数据变更后不再满足 INNER JOIN 或过滤条件,历史投影行不会被自动撤回。
  • 没有跨源一致快照:多个增量分支读取的是构建时各源表的当前数据,不保证得到同一时刻的跨源快照。
  • 分支数量会影响成本:每张增量表对应一个完整视图分支,增量表越多,源表扫描和 UNION DISTINCT 去重开销越大。
  • 水位来自投影结果:无法产生最终关联结果的较大批次值不会推进水位,因此可能在后续构建中被重复扫描。这可以保证关联数据迟到后仍能补齐结果。
  • INSERT 不保证幂等:对重复写入敏感的业务,建议选择 MERGE 并使用稳定的业务唯一条件。

10. 验证构建结果

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

  1. 确认首次构建为全量构建,后续构建为增量构建。
  2. 检查任务状态是否成功,以及输出数据量是否符合预期。
  3. 在 SQL 优化信息或 SQL 翻译详情中,检查每张增量表是否分别使用了 > 上次水位 条件。
  4. 多表场景检查各分支是否使用 UNION 合并,而不是将多个增量条件同时放入一个分支。
  5. 根据投影配置,确认最终写入使用 INSERT INTOMERGE INTO
  6. MERGE 场景,校验业务唯一键是否只保留一条结果。

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