订单1001先创建,随后付款金额从 100 元改成 80 元,物流系统又补来运单号;一条版本更旧的取消消息最后才到。四张 Paimon 表使用相同主键,却可能分别返回最新整行、拼接宽行、累计金额和首次事件。
问题不在 Paimon 算出了四个互相矛盾的答案,而在建表人没有先说明:输入究竟是完整快照、字段补丁、可累加增量,还是只用于判重的首次事件。
merge-engine不是性能开关。主键只圈定哪些记录需要冲突处理,Merge Engine 决定这些记录最终代表覆盖、补列、聚合,还是只认第一次。
同一个订单事实,必须先翻译成四种输入契约
先固定共同业务事实和到达顺序,机器规模、吞吐与 SLA 均未提供,因此本文只验证语义,不比较性能。
| 到达次序 | 业务版本 | 事实 |
|---|---|---|
| 1 | 1 | 创建订单,应付金额 100 元 |
| 2 | 3 | 支付成功,应付金额调整为 80 元 |
| 3 | 4 | 物流系统补充运单号SF001 |
| 4 | 2 | 旧的取消事件迟到 |
这四条事实不能不加转换地喂给所有 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)。注意第三条必须重复携带status和amount;若物流源只发送tracking_no,deduplicate会保留那条稀疏记录,而不会自动从旧行补列。
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 只补运单号,status与amount沿用此前非空值;版本 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_delta的sum对输入顺序不敏感,max_version记录目前观察到的最大版本。没有配置聚合函数的非主键字段默认使用last_non_null_value,这经常让一张表同时混合可逆与不可逆语义,审查时不能只看merge-engine=aggregation。
Aggregation 文档 还给出更严格的撤回边界:并非所有聚合函数都支持UPDATE_BEFORE与DELETE。即使某些函数允许忽略撤回,也等于主动改变结果语义,不能拿任务持续运行代替业务对账。
first-row:认的是合并顺序中的第一次,不是最小业务版本
first-row用于只保留同主键第一次出现,例如设备首次激活或日志去重。它不能配置sequence.field,也默认不接受DELETE与UPDATE_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_deduplicate | PAID, 80, SF001, v4 | 最大 Sequence 的完整行获胜 | 稀疏更新会自动补齐、下游 Changelog 完整 |
order_partial_update | PAID, 80, SF001, v4 | 非空字段按版本拼接 | NULL 可以清空字段、Delete 已被正确传播 |
order_aggregation | 80, max(v)=4 | 输入增量按声明函数结合 | 80 是当前行快照、所有聚合都支持撤回 |
order_first_row | CREATED, 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.0,CoreOptions.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。配置选择DeduplicateMergeFunction,sequence.field改变同键KeyValue进入它的顺序。其余三种引擎应复制这套程序并替换 DDL 与输入契约,不能只换参数而复用错误数据。
建表评审先问四句话,再看参数
选择 Merge Engine 前,把下面四句话写进数据契约:
- 同一主键的每条消息是完整状态、字段补丁,还是可结合的增量?
- 新旧顺序由哪个业务版本决定,它是否跨来源可比较?
- NULL 表示未提供还是主动清空,Delete 是删整行、撤回部分字段还是应被拒绝?
- 下游只查当前状态,还是还要消费完整撤回流?
对应选择可以很直接:
- 完整行快照、最大版本覆盖:
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
DeduplicateMergeFunctionPartialUpdateMergeFunctionAggregateMergeFunctionFirstRowMergeFunction