Apache Paimon 实时湖仓实战(第 2 篇):同一批订单写进四种 Merge Engine,为什么得到四个答案
2026/9/7 11:18:47 网站建设 项目流程

订单1001先创建,随后付款金额从 100 元改成 80 元,物流系统又补来运单号;一条版本更旧的取消消息最后才到。四张 Paimon 表使用相同主键,却可能分别返回最新整行、拼接宽行、累计金额和首次事件。

问题不在 Paimon 算出了四个互相矛盾的答案,而在建表人没有先说明:输入究竟是完整快照、字段补丁、可累加增量,还是只用于判重的首次事件。

merge-engine不是性能开关。主键只圈定哪些记录需要冲突处理,Merge Engine 决定这些记录最终代表覆盖、补列、聚合,还是只认第一次。

同一个订单事实,必须先翻译成四种输入契约

先固定共同业务事实和到达顺序,机器规模、吞吐与 SLA 均未提供,因此本文只验证语义,不比较性能。

到达次序业务版本事实
11创建订单,应付金额 100 元
23支付成功,应付金额调整为 80 元
34物流系统补充运单号SF001
42旧的取消事件迟到

这四条事实不能不加转换地喂给所有 Merge Engine。不同引擎要求的是不同数据产品:

Merge Engine表要回答的问题每条输入必须代表什么正确结果
deduplicate订单最新完整状态是什么一份可以独立解释的完整行快照最大业务版本对应的整行
partial-update多个来源怎样补成当前宽行只携带本次需要更新的列各字段按自己的更新顺序拼接
aggregation订单累计发生了多少金额变化可按指定函数结合的增量或状态各字段的聚合结果
first-row这个订单第一次出现时是什么只追加的候选首次事件物理合并顺序中的第一行

因此,公平对比不是让四张表吃完全相同的列值,而是让它们面对同一业务历史,再按各自声明的输入契约投影数据。如果上游发的是整行 CDC,partial-update未必需要;如果金额是当前余额,直接交给sum就会重复累计。

四张表的 DDL 已经写下了四种业务真相

下面的最小实验固定为 Paimon 2.0.0 与 Flink 1.20,本地文件系统只用于隔离语义。所有表都使用order_id作为主键、4 个固定 Bucket。先关闭 Flink Sink 的 Upsert Materialize,避免它在写入 Paimon 前再次改写乱序语义;这是 Paimon Merge Engine 文档 对 Flink SQL 的明确要求。

SET'table.exec.sink.upsert-materialize'='NONE';

deduplicate:保留最后合并的完整行

deduplicate是默认模式。配置sequence.field=source_version后,最大版本最后参与合并;版本相同才退回输入顺序。

CREATETABLEorder_deduplicate(order_idBIGINT,statusSTRING,amountDECIMAL(18,2),tracking_no STRING,source_versionBIGINT,PRIMARYKEY(order_id)NOTENFORCED)WITH('bucket'='4','merge-engine'='deduplicate','sequence.field'='source_version');INSERTINTOorder_deduplicateVALUES(1001,'CREATED',100.00,CAST(NULLASSTRING),1),(1001,'PAID',80.00,CAST(NULLASSTRING),3),(1001,'PAID',80.00,'SF001',4),(1001,'CANCELLED',100.00,CAST(NULLASSTRING),2);

预期结果是版本 4 的完整行:(1001, PAID, 80.00, SF001, 4)。注意第三条必须重复携带statusamount;若物流源只发送tracking_nodeduplicate会保留那条稀疏记录,而不会自动从旧行补列。

DeduplicateMergeFunction的核心状态只有latestKv:每加入一条非忽略记录,就让它替换前一条。Sequence 决定加入顺序,Merge Function 负责留下最后一条。

partial-update:NULL 默认表示这次不更新

物流、支付、订单三个来源各自只拥有部分字段时,partial-update才是对应语义。它逐列使用同主键下的最新数据,默认不让 NULL 覆盖旧值。

CREATETABLEorder_partial_update(order_idBIGINT,statusSTRING,amountDECIMAL(18,2),tracking_no STRING,source_versionBIGINT,PRIMARYKEY(order_id)NOTENFORCED)WITH('bucket'='4','merge-engine'='partial-update','sequence.field'='source_version');INSERTINTOorder_partial_updateVALUES(1001,'CREATED',100.00,CAST(NULLASSTRING),1),(1001,'PAID',80.00,CAST(NULLASSTRING),3),(1001,CAST(NULLASSTRING),CAST(NULLASDECIMAL(18,2)),'SF001',4),(1001,'CANCELLED',100.00,CAST(NULLASSTRING),2);

预期结果同样是(1001, PAID, 80.00, SF001, 4),但形成路径完全不同:版本 4 只补运单号,statusamount沿用此前非空值;版本 2 因 Sequence 更小,不能把状态改回取消。

这里最危险的不是迟到,而是 NULL 的含义。默认规则把 NULL 当作未提供,因此不能表达把tracking_no主动清空。多条输入流各自维护不同字段时,一个全局 Sequence 还可能被另一条流覆盖;Partial Update 文档 为此提供 Sequence Group,让每组字段拥有自己的版本边界。

Delete 也不是默认可接受输入。Paimon 2.0.0 要求显式选择:忽略删除、收到删除时移除整行,或使用 Sequence Group 撤回部分字段。没定义删除契约就把 CDC 直灌partial-update,失败或静默丢失业务删除都不应被称为正确。

aggregation:80 是累计出来的,不是覆盖出来的

如果表要保存的是订单金额净变化,上游必须投影为增量:创建+100,金额调整-20,状态与物流事件不贡献金额。把两份完整金额 100 和 80 直接求和得到 180,是输入契约错误,不是聚合误差。

CREATETABLEorder_aggregation(order_idBIGINT,amount_deltaDECIMAL(18,2),max_versionBIGINT,PRIMARYKEY(order_id)NOTENFORCED)WITH('bucket'='4','merge-engine'='aggregation','fields.amount_delta.aggregate-function'='sum','fields.max_version.aggregate-function'='max');INSERTINTOorder_aggregationVALUES(1001,100.00,1),(1001,-20.00,3),(1001,0.00,4),(1001,0.00,2);

预期结果是(1001, 80.00, 4)amount_deltasum对输入顺序不敏感,max_version记录目前观察到的最大版本。没有配置聚合函数的非主键字段默认使用last_non_null_value,这经常让一张表同时混合可逆与不可逆语义,审查时不能只看merge-engine=aggregation

Aggregation 文档 还给出更严格的撤回边界:并非所有聚合函数都支持UPDATE_BEFOREDELETE。即使某些函数允许忽略撤回,也等于主动改变结果语义,不能拿任务持续运行代替业务对账。

first-row:认的是合并顺序中的第一次,不是最小业务版本

first-row用于只保留同主键第一次出现,例如设备首次激活或日志去重。它不能配置sequence.field,也默认不接受DELETEUPDATE_BEFORE

CREATETABLEorder_first_row(order_idBIGINT,first_status STRING,first_amountDECIMAL(18,2),observed_versionBIGINT,PRIMARYKEY(order_id)NOTENFORCED)WITH('bucket'='4','merge-engine'='first-row');INSERTINTOorder_first_rowVALUES(1001,'CREATED',100.00,1),(1001,'PAID',80.00,3),(1001,'CANCELLED',100.00,2);

在这组顺序稳定的追加输入中,预期结果是(1001, CREATED, 100.00, 1)。但若版本 3 先到,Paimon 不会根据observed_version回头选择版本 1。它保证的是同主键合并后只留第一行,不是按业务时间求最早。

First Row 文档 还有一个容易漏掉的可见性代价:L0 文件要经过 Compaction 才可见,因此默认同步 Compaction;切成异步后可能增加数据可见延迟。它产生 Insert-Only Changelog 的优势,来自更严格的输入限制,并非免费获得。

一条 SELECT 只能证明最终表状态

四组写入完成后,分别执行只读查询:

SELECT*FROMorder_deduplicateWHEREorder_id=1001;SELECT*FROMorder_partial_updateWHEREorder_id=1001;SELECT*FROMorder_aggregationWHEREorder_id=1001;SELECT*FROMorder_first_rowWHEREorder_id=1001;

预期证据矩阵如下:

关键结果结果支持什么不能证明什么
order_deduplicatePAID, 80, SF001, v4最大 Sequence 的完整行获胜稀疏更新会自动补齐、下游 Changelog 完整
order_partial_updatePAID, 80, SF001, v4非空字段按版本拼接NULL 可以清空字段、Delete 已被正确传播
order_aggregation80, max(v)=4输入增量按声明函数结合80 是当前行快照、所有聚合都支持撤回
order_first_rowCREATED, 100, v1当前到达顺序下首行被保留乱序时仍能得到最小业务版本

再查询 Snapshot,只能确认这些写入形成了可见提交:

-- 观察对象:每张表的提交时间线;只读。-- 正常信号:写入对应的 Snapshot 可见,commit_kind 符合预期。-- 异常分支:无新 Snapshot 时先查写入与 Commit;有 Snapshot 但值错时回查输入契约、Sequence 与 Merge Engine。SELECTsnapshot_id,commit_kind,commit_time,total_record_count,changelog_record_countFROMorder_partial_update$snapshotsORDERBYsnapshot_id;

这组实验能证明四种表内合并语义不同,也能暴露完整行、字段补丁、指标增量和首次事件的契约差异。它不能证明生产吞吐、Compaction 成本、故障恢复连续性,也不能证明流式下游收到完整UPDATE_BEFORE/UPDATE_AFTER;后者属于第 03 篇的 Changelog Producer 边界。

源码里没有“猜业务”,只有四套确定的状态机

固定到release-2.0.0CoreOptions.MergeEngine定义四个枚举值,默认是DEDUPLICATE。建表配置最终让 Primary Key Table 选择对应的 Merge Function:

同主键 KeyValue 按 Sequence / 内部顺序进入合并 ├─ DeduplicateMergeFunction:不断覆盖 latestKv ├─ PartialUpdateMergeFunction:按字段、Sequence Group 与 NULL 规则更新 ├─ AggregateMergeFunction:为各字段调用声明的 FieldAggregator └─ FirstRowMergeFunction:first 为空时赋值,之后不再替换

源码能证明的是状态如何变化,不能替业务决定 NULL 是清空还是缺失、金额是余额还是增量、第一次按到达时间还是业务时间定义。配置能够忠实执行错误契约,这恰恰是 Merge Engine 最危险的地方。

Java 用相同输入断言 Merge Engine 结果

以下程序按paimon-flink-1.20:2.0.0API 编写,args[0]是隔离测试 Warehouse;本环境未启动 Flink/Paimon 集群运行全文示例。

importorg.apache.flink.table.api.EnvironmentSettings;importorg.apache.flink.table.api.TableEnvironment;publicfinalclassDeduplicateContractCheck{publicstaticvoidmain(String[]args)throwsException{if(args.length!=1)thrownewIllegalArgumentException("warehouse is required");TableEnvironmentt=TableEnvironment.create(EnvironmentSettings.newInstance().inBatchMode().build());t.executeSql("CREATE CATALOG p WITH ('type'='paimon','warehouse'='"+args[0]+"')");t.executeSql("USE CATALOG p");t.executeSql("CREATE DATABASE IF NOT EXISTS demo");t.executeSql("CREATE TABLE demo.order_dedup (id BIGINT,status STRING,ver BIGINT,"+"PRIMARY KEY(id) NOT ENFORCED) WITH "+"('merge-engine'='deduplicate','sequence.field'='ver','bucket'='4')");t.executeSql("INSERT INTO demo.order_dedup VALUES "+"(1001,'PAID',3),(1001,'CANCELLED',2)").await();t.executeSql("SELECT id,status,ver FROM demo.order_dedup WHERE id=1001").print();}}

预期只输出PAID,3;若输出旧状态,先查 Sequence 输入,不应先调 Compaction。配置选择DeduplicateMergeFunctionsequence.field改变同键KeyValue进入它的顺序。其余三种引擎应复制这套程序并替换 DDL 与输入契约,不能只换参数而复用错误数据。

建表评审先问四句话,再看参数

选择 Merge Engine 前,把下面四句话写进数据契约:

  1. 同一主键的每条消息是完整状态、字段补丁,还是可结合的增量?
  2. 新旧顺序由哪个业务版本决定,它是否跨来源可比较?
  3. NULL 表示未提供还是主动清空,Delete 是删整行、撤回部分字段还是应被拒绝?
  4. 下游只查当前状态,还是还要消费完整撤回流?

对应选择可以很直接:

  • 完整行快照、最大版本覆盖:deduplicate + sequence.field
  • 多来源补列:partial-update,并评估 Sequence Group 与 Delete 规则;
  • 输入天然是增量且聚合函数符合撤回要求:aggregation
  • 输入只追加、到达第一次就是业务第一次:first-row

任何一句答不清,都不应先创建生产表。尤其不要把partial-update当作节省上游补全成本的万能模式,也不要把aggregation当作写入时顺手预计算。

选错 Merge Engine 的后果不是查询慢一点,而是所有 Snapshot 都稳定保存了错误的业务世界。

面试表达主线

Paimon 的主键只标识冲突集合,Merge Engine 才定义冲突如何收敛:deduplicate保留最后完整行,partial-update按列补齐,aggregation按字段函数结合,first-row保留第一次到达。选型先定义输入、顺序、NULL、Delete 和 Changelog 契约;否则技术提交成功也可能持续产出错误业务状态。

官方资料

  • Merge Engine
  • Partial Update
  • Aggregation
  • First Row
  • Sequence Field and RowKind
  • DeduplicateMergeFunction
  • PartialUpdateMergeFunction
  • AggregateMergeFunction
  • FirstRowMergeFunction

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询