:Paimon 表里明明是 80 元,下游为何算成 180?Checkpoint 成功不等于变更完整)
先把问题放进一个常见的电商结算场景订单原价 100 元客服改价后变成 80 元。Paimon 明细表查询正确Flink Checkpoint 也连续成功可下游按门店汇总的销售额却从 100 加到了 180。如果只盯着“表里是 80”和“Checkpoint 是绿的”这起事故看起来毫无道理。真正缺失的不是 80 元这个新值而是“先把旧的 100 元减掉”那条撤回消息。Exactly-Once 只能保证同一提交不会多算一次不能凭空补出下游聚合所需的UPDATE_BEFORE。下文的配置、源码和行为判断固定在Apache Paimon 2.0.0源码基线为release-2.0.0/604e6d5e...。先找到丢失的那一步80 元到了100 元没有撤回Changelog Producer 官方文档 明确区分四种模式模式旧值从哪里来主要代价适用前提none不额外保存消费者可能自行 Normalize下游状态可能很大下游只需 Upsert或能按主键记住旧值input直接保存上游完整 Changelog额外 Changelog 文件输入确实含完整 Before/AfterlookupCompaction 时查旧值生成本地缓存、磁盘与 Compaction 成本输入不完整但下游必须撤回full-compaction两次全量合并结果做 Diff高写放大与更高延迟能接受较大延迟和全量合并成本默认none下Paimon 能表达某个主键的新值却不一定直接给出旧值。下游若按store_id聚合而表主键是order_id没有旧金额就不知道该先减 100 再加 80。为什么 Checkpoint 成功仍然救不了这笔汇总源端成功CDC 事件完整到达 提交成功Checkpoint 对应的 Commit Message 进入 Snapshot 消费正确下游收到足以更新自己状态的 RowKind 序列最小实验不要只数行。写入I(100)再更新为-U(100), U(80)分别观察 Paimon 流读输出和下游聚合。若上游只给U(80)却配置inputPaimon 不会发明缺失的-U(100)。CREATETABLEorders(order_idBIGINT,store_idBIGINT,amountDECIMAL(18,2),PRIMARYKEY(order_id)NOTENFORCED)WITH(changelog-producerlookup);这条 DDL 不是默认建议。lookup把旧值查询放进 Compaction 路径会使用本地内存与磁盘缓存并可能拉长提交。官方还特别提醒 Changelog Producer 可能显著降低 Compaction 性能。别再只数行同一笔更新必须核对 RowKind把同一订单写入四张仅修改changelog-producer的表才能公平验证差异。共同前提是表主键为order_id下游按store_id汇总输入先为I(1001,10,100)再把金额改成 80。正确撤回流I(1001,10,100) → -U(1001,10,100) → U(1001,10,80) 危险更新流I(1001,10,100) → U(1001,10,80)第二条对主键 Upsert Sink 可能足够对SUM(amount) GROUP BY store_id却缺少减 100 的证据。若结果变成 180先查审计 Sink 中的 RowKind再查 Checkpoint顺序反过来会把语义缺失误诊成重复提交。固定到release-2.0.0旧值生成和读取的最短源码链是MergeTree Compaction → LookupChangelogMergeFunctionWrapper / FullChangelogMergeFunctionWrapper → ChangelogResult → Changelog File → IncrementalChangelogReadProvider源码能证明变化在哪里补齐和落盘不能证明缓存容量足够也不能证明消费者按 RowKind 正确处理。-- 观察提交是否携带 Changelog只读。-- 有记录不能替代逐主键 RowKind 与业务聚合对账。SELECTsnapshot_id,commit_kind,total_record_count,changelog_record_countFROMorders$snapshotsORDERBYsnapshot_id;本实验能证明 Changelog 契约是否足够不能证明生产规模下的性能。机器、Bucket、更新率与 SLA 未提供必须另外压测 Checkpoint P99、Compaction 时长、缓存占用和 Changelog 文件增速。该选哪种模式答案藏在下游怎么消费数据库型 Sink 能按主键覆盖时none往往已经够用非主键聚合需要撤回旧值应先确认输入是否自带完整 Before/After再在input与lookup之间选择。full-compaction适合能够接受更长生成周期的场景不能仅为省掉 Normalize State 就默认打开。用 Java 把 RowKind 打出来别再猜撤回有没有发生固定版本源码如何闭环到输出源码固定为 Apache Paimonrelease-2.0.0/604e6d5e...paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java的getResult()在包含 L0 且策略要求 Changelog 时调用setChangelog(before, after)决定是否产生-U/Upaimon-core/src/main/java/org/apache/paimon/mergetree/compact/FullChangelogMergeFunctionWrapper.java的getResult()比较合并前后的记录决定完整 Compaction Changelog对应源码Lookup wrapper、Full wrapper。Java 中的Row.getKind()正是最终观测点若金额 100→80 只有U(80)而没有旧值撤回下游非主键聚合就缺少减 100 的依据。它证明输出契约不证明消费者已经正确执行撤回。以下程序按paimon-flink-1.20:2.0.0API 编写args[0]指向已有测试表的 Warehouse流式collect()会持续读取生产应接入可持久化审计 Sink。本环境未运行全文示例。importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;importorg.apache.flink.table.api.EnvironmentSettings;importorg.apache.flink.table.api.TableResult;importorg.apache.flink.table.api.bridge.java.StreamTableEnvironment;importorg.apache.flink.types.Row;importorg.apache.flink.util.CloseableIterator;publicfinalclassChangelogAuditJob{publicstaticvoidmain(String[]args)throwsException{if(args.length!1)thrownewIllegalArgumentException(warehouse is required);Stringwarehouseargs[0].replace(,);StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();env.enableCheckpointing(30_000L);StreamTableEnvironmenttStreamTableEnvironment.create(env,EnvironmentSettings.newInstance().inStreamingMode().build());t.executeSql(CREATE CATALOG p WITH (typepaimon,warehousewarehouse));t.executeSql(USE CATALOG p);TableResultresultt.executeSql(SELECT order_id,amount FROM demo.orders);try(CloseableIteratorRowrowsresult.collect()){while(rows.hasNext()){Rowrowrows.next();System.out.printf(%s id%s amount%s%n,row.getKind(),row.getField(0),row.getField(1));}}}}生产中应把 Print 换成可持久化审计 Sink。输出缺少旧金额的-U时Checkpoint 成功不能挽救非主键聚合。lookup/full-compaction分别进入LookupChangelogMergeFunctionWrapper或FullChangelogMergeFunctionWrapper随后由IncrementalChangelogReadProvider供流读。数据库型 Sink 通常能按主键 Upsertnone往往足够Flink 非主键聚合需要完整撤回应该在上游已完整时用input否则评估lookup能容忍较长延迟、又必须从表状态推导完整变化时再考虑full-compaction。不要为了消除 Flink Normalize 就无条件打开lookup。应同时压测Checkpoint P95/P99、Compaction 时长、本地缓存占用、Changelog 文件数以及故障恢复后的业务聚合对账。Paimon 表里的最终值正确只证明状态存对了下游能不能算对还取决于变化过程有没有被完整表达。官方资料Changelog ProducerPrimary Key TableSnapshot SpecificationFlink Consumer ID