ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Apache Paimon 数据出仓源码导读(二):Paimon Changelog 到底输出什么?+I、-U、+U、-D 与 Changelog Producer

Apache Paimon 数据出仓源码导读(二):Paimon Changelog 到底输出什么?+I、-U、+U、-D 与 Changelog Producer 上一篇我们把 Flink 从 Paimon 读取“全量加增量”的接力过程走了一遍先固定最新 Snapshot 做全量再从下一个 Snapshot 继续追增量。任务跑起来以后新的疑问很快就来了。同样是一条订单 UPDATE为什么有时看到I[1001, 100.00, CREATED]有时是-D[1001, 80.00, CREATED] I[1001, 100.00, CREATED]还有时又变成-U[1001, 80.00, CREATED] U[1001, 100.00, CREATED]三种结果看起来完全不同到底哪一种才算正确答案是它们都可能正确只是记录变化的方式不同。RowKind说明“这一行要怎样作用到结果表上”changelog-producer决定“Paimon 从哪里拿到这些变化以及有没有能力还原修改前的旧值”。这一篇我们继续跟着订单1001看看它从80元改成100元以后为什么会走出三种不同的 Changelog。一、今天的现场同一条 UPDATE打印出来却不一样订单表还是上一篇的ordersCREATETABLEorders(idBIGINT,amountDECIMAL(10,2),statusSTRING,PRIMARYKEY(id)NOTENFORCED)WITH(connectorpaimon,pathhdfs:///warehouse/orders,changelog-producerinput);在 SnapshotS10中订单1001是id1001, amount80.00, statusCREATED随后 MySQL 执行UPDATEordersSETamount100.00WHEREid1001;在业务眼里这只是一条 UPDATE。但数据系统要把它交给下游时至少要回答两个问题修改前的80.00还要不要发出去如果要发应该把旧行标成-D还是-U这两个问题并不是由 SQL 中的UPDATE三个字直接决定而是由上游传进来的 RowKind、Paimon 的 Changelog Producer以及 Flink 下游需要什么样的 Changelog 一起决定。先别急着看 Producer我们先把四个符号说成人话。二、I、-U、U、-D其实是四张“记账凭证”Paimon 的RowKind.java定义了四种行变化INSERT(I,(byte)0),UPDATE_BEFORE(-U,(byte)1),UPDATE_AFTER(U,(byte)2),DELETE(-D,(byte)3);不要把它们理解成四种文件也不要把和-当成正负数。它们表达的是这条记录到达下游后应该加入结果还是从结果中撤回RowKind大白话对结果表做什么常见位置I以前没有现在新增一行加入新记录INSERT或者按主键覆盖的新值-U这是一行“修改前的旧值”先撤回旧记录UPDATE 的前半段U这是一行“修改后的新值”再加入新记录UPDATE 的后半段-D这行已经不存在了删除原记录DELETE源码还把它们分成了两组publicbooleanisRetract(){returnthisUPDATE_BEFORE||thisDELETE;}publicbooleanisAdd(){returnthisINSERT||thisUPDATE_AFTER;}也就是说-U和-D都属于“撤回”I和U都属于“加入”。这个分组比符号的名字更重要。很多下游算子真正关心的是“加一条还是减一条”而不是它在业务 SQL 中叫 INSERT 还是 UPDATE。三、-U和-D都在撤回为什么还要分两个名字我们先看订单被真正删除-D[1001, 80.00, CREATED]它的意思是订单 1001 到这里就结束了后面没有替代它的新版本。再看订单金额从80改成100-U[1001, 80.00, CREATED] U[1001, 100.00, CREATED]这里的-U不是说订单被业务删除了而是在告诉下游先把旧版本撤回马上会有一个新版本接替它。如果下游只是一个按id覆盖的主键表它只收到新值也能工作找到id1001把80覆盖成100就行。但如果下游正在算订单总金额旧值就不能省略原合计80 -U 旧值减去 80 U 新值加上 100 新合计100没有-U[80]下游只看到U[100]就不知道应该减掉多少。除非它自己按主键保存每一行的旧值。这正是完整 Changelog 有价值的地方它把“查旧值”的工作提前做了。四、先分清全量和增量第一次读到的通常都是I上一篇讲过scan.modelatest-full会先读取启动时最新 Snapshot 的完整表状态。假设启动时固定的是S10订单1001在S10中已经存在。那么第一阶段输出的是I[1001, 80.00, CREATED]原因很简单对一个刚启动、结果为空的下游来说S10中现存的每一行都是“加入基线”。下游并不需要知道它在过去被改过多少次。到了S11金额从80变成100Producer 的差异才真正出现。所以排查时要先问一句我现在看到的是第一次全量基线还是后续 Snapshot 的增量变化如果把两段输出混在一起很容易误以为 Producer 没生效。五、changelog-producer决定的不是符号长相而是“旧值从哪里来”CoreOptions.ChangelogProducer一共定义了四种模式publicenumChangelogProducer{NONE(none,No changelog file.),INPUT(input,... changelog is from input.),FULL_COMPACTION(full-compaction,Generate changelog files with each full compaction.),LOOKUP(lookup,Generate changelog files through lookup compaction.);}默认值是changelog-producernone四种模式真正回答的是四个不同的问题Producer旧值从哪里来是否写 Changelog 文件核心特点none不找旧值否直接读取新增数据文件成本最低input上游传什么就保存什么是最接近原始输入保留中间变化lookup按主键查询历史状态是可以还原逻辑上的 before/afterfull-compaction全量 Compaction 时比较旧状态和新状态是输出完整变化但有批次延迟它们不是“越来越高级”的四个档位而是把成本放在了不同位置。六、同一条订单 UPDATE在四种模式下会看到什么现在把案例固定下来S101001, 80.00, CREATED S111001, 100.00, CREATED而且这张表由 Paimon MySQL CDC Action 写入。第六篇入仓源码导读中讲过MySqlRecordParser.extractRecords()会把 before 写成 DELETE把 after 写成 INSERTif(!before.isEmpty()){records.add(createRecord(RowKind.DELETE,before));}if(!after.isEmpty()){records.add(createRecord(RowKind.INSERT,after));}于是同一个业务 UPDATE 可能得到下面四种结果Producer增量阶段可能输出下游应该怎样理解noneI[1001, 100]按主键把新值写进去旧值不单独输出input-D[1001, 80]、I[1001, 100]原样保留 MySQL CDC Action 的 before/after 表达lookup-U[1001, 80]、U[1001, 100]Paimon 查到旧状态后重新判断为逻辑 UPDATEfull-compaction-U[1001, 80]、U[1001, 100]下一次 Full Compaction 时输出这段净变化这里需要补一句none输出的新值也可能是U取决于写入端原本传入的 RowKind。对none来说最重要的保证不是“更新一定长成I”而是它提供的是 Upsert Changelog不保证给出-U旧值。另外input名字虽然叫 Changelog Producer但它不会把-D/I自动改名成-U/U。它的原则恰恰是输入是什么就保存什么。七、none不单独记流水账只把新文件交给下游当 Producer 是none时DataTableStreamScan.createFollowUpScanner()选择caseNONE:followUpScannernewDeltaFollowUpScanner();break;DeltaFollowUpScanner只处理APPENDSnapshot并使用ScanMode.DELTA它读取的是本次 Snapshot 新增的数据文件而不是一份单独保存的 Changelog 文件。主键表写入数据文件前同一个 Key 在 Write Buffer 中会经过 Merge Function 合并。默认的DeduplicateMergeFunction只保留最新一条publicvoidadd(KeyValuekv){latestKvkv;}因此MySQL UPDATE 写进来的-D[1001, 80] I[1001, 100]如果在同一次缓冲区合并中处理数据文件留下的最终结果可能只有I[1001, 100]旧行已经在合并过程中消失了Paimon Source 无法再凭空生成-U[80]。为什么默认值要设成none因为不是所有下游都需要旧值。如果下游是一个带主键的 Upsert Kafka、数据库或 OLAP 表收到id1001的新值后直接覆盖即可。为了一个用不到的-U让每次写入都多写 Changelog 文件或者查询历史状态代价并不划算。BaseDataTableSource.getChangelogMode()对这类流式主键表返回returnChangelogMode.upsert();Upsert Mode 通常包含I、U、-D不包含-U。这是一份“按主键应用就能得到正确最终状态”的变化流不是一份完整审计流水。八、没有 Producer为什么最后的 Flink 结果里仍可能出现-U这是最容易让人怀疑人生的一幕表明明配置了changelog-producernonePrint Sink 却打印出了-U/U。原因可能不在 Paimon而在 Flink Planner。Paimon Source 告诉 Flink我能提供 Upsert Changelog并且这张表有主键 id。如果下游算子必须拿到完整的 before/afterFlink 可以在执行计划中加入 Changelog Normalize。它按主键保存上一版记录PaimonI[1001, 100] Normalize 状态中已有[1001, 80] Flink 向下游补成 -U[1001, 80] U[1001, 100]Paimon 仓库中的随机集成测试也直接写着// changelog is produced by Flink normalize operator所以“最终 Print Sink 打印了什么”和“Paimon 原生 Changelog 文件里有什么”不是同一个问题。这种设计把压力从表的写入侧挪到了消费作业位置付出的代价Paimon Producer 生成旧值写入、Compaction、Changelog 文件存储成本Flink Normalize 生成旧值下游状态、Checkpoint、恢复时间成本排查时可以先执行EXPLAIN看看计划中是否出现了 Changelog Normalize 一类的节点。九、input一边写最终数据一边保存上游原始流水input的源码入口非常直白。MergeTreeWriter.flushWriteBuffer()只有在 INPUT 模式下才创建 Changelog WriterfinalRollingFileWriterKeyValue,DataFileMetachangelogWriterchangelogProducerChangelogProducer.INPUT?writerFactory.createRollingChangelogFileWriter(0):null;随后SortBufferWriteBuffer.forEach()同时接收两个 ConsumerrawConsumer记录合并前的原始输入mergedConsumer记录同 Key 合并后的最终结果。源码中的读取顺序是先把每条原始 KeyValue 交给rawConsumer再把 Merge Function 的结果交给数据文件 Writer。这就是 CoreOptions 描述中的 “Double write”。同一批输入走出了两本账数据文件服务于“这张表现在是什么样” Changelog 文件服务于“上游刚才传来了哪些行变化”input最重要的优点它能保留中间变化。如果一个 Flush 周期内订单金额连续变化80 → 90 → 100数据文件最终只需要留下100但原始 Changelog 仍可以保留上游送来的多次变化。这对审计、规则触发和需要逐笔事件的场景很有价值。不过“保留每次行变化”不等于“保留 MySQL Binlog 的全局事务顺序”。Write Buffer 落盘前仍会排序多 Bucket、多并行度之间也没有一条统一的事件时间线。input保存的是表级 Changelog不是 Binlog 归档。如果业务需要严格按照数据库事务顺序回放还是应该保留原始 Binlog 或 Kafka CDC 流。input最容易被误解的地方它只负责“原样记账”不负责“补齐旧值”。上游传-U/U它就保存-U/UMySQL CDC Action 传-D/I它就保存-D/I上游只有新值它也变不出旧值。代价则是写放大正常数据文件要写Changelog 文件也要写Manifest 和后续清理也会多一份工作。十、lookup先翻旧账再判断这是新增、更新还是删除lookup的目标不是照抄输入而是回答一个逻辑问题这个 Key 在变化前存在吗变化后还存在吗前后内容一样吗核心代码在LookupChangelogMergeFunctionWrapper.getResult()。它先尝试从参与 Compaction 的高层文件中找到旧记录找不到时再按 Key 去历史数据中 Lookupif(highLevelnull){TlookupResultlookup.apply(mergeFunction.key());...}拿到 before再把 Level 0 中的新变化合并成 after最后进入setChangelog(before, after)if(beforenull||!before.isAdd()){if(after.isAdd()){add(INSERT);}}else{if(!after.isAdd()){add(DELETE);}elseif(valueChanged){add(UPDATE_BEFORE);add(UPDATE_AFTER);}}换成大白话就是beforeafter输出不存在存在I存在不存在-D内容使用 before存在存在且变化-U before、U after都不存在无输出这也是为什么 lookup 能把 MySQL CDC Action 的-D/I重新整理成更标准的逻辑 UPDATE-U[1001, 80] U[1001, 100]为什么不让所有表都默认 lookup因为“翻旧账”不是免费的。它需要维护或读取 Lookup 状态还要让 Level 0 变化参与 Compaction。lookup-wait默认是truekey(lookup-wait).defaultValue(true)源码中的含义是需要 Lookup 时Commit 默认会等待 Lookup Compaction。换来的好处是下游较快拿到完整 before/after付出的代价是写入和提交链路更重热点 Key、高基数主键、大 Bucket 都可能放大 Lookup 压力。十一、full-compaction不逐笔查旧值等大盘点时一次算清full-compaction也能产生完整 Changelog但思路不同。它不在每次变化到来时按 Key Lookup而是在 Full Compaction 中比较上一次进入最高层的稳定记录 vs 本次合并后的最终记录FullChangelogMergeFunctionWrapper中topLevelKv代表上一次 Full Compaction 后的旧状态merged代表本轮把增量全部合并后的新状态。比较规则与 lookup 类似旧状态没有新状态有-I旧状态有新状态没有--D旧状态有新状态有且变化--U、U但FullChangelogMergeTreeCompactRewriter只有在输出到最高 Level 时才生成 ChangelogbooleanchangelogoutputLevelmaxLevel;所以它最大的特点是“批量结账”。假设两次 Full Compaction 之间发生了80 → 90 → 100最终更可能看到的是-U[80] U[100]中间的90已经被合并掉了。这不适合逐笔审计却很适合只关心两个稳定版本之间净变化的场景。Flink 写入链路可以通过下面的配置控制 Full Compaction 触发节奏full-compaction.delta-commits changelog-producer.compaction-interval间隔越大Compaction 次数越少但 Changelog 延迟越高间隔越小变化更快可见但 Full Compaction 压力更频繁。十二、为什么lookup和full-compaction可能忽略“中间过程”很多人看到lookup能生成-U/U就以为它等同于数据库 Binlog能逐条还原每一次业务 UPDATE。并不一定。LookupChangelogMergeFunctionWrapper和FullChangelogMergeFunctionWrapper最终比较的都是一个 Compaction 窗口开始前的状态 一个 Compaction 窗口合并后的状态如果同一个 Key 在这个窗口内连续更新多次Merge Function 会先得到最终结果Changelog 再比较 before 和 after。中间版本可能已经合并。四种 Producer 的“记忆粒度”可以这样理解Producer更像什么是否适合逐笔回放none只告诉你新的 Upsert 结果否input复印上游交来的原始流水最适合lookup每次 Lookup Compaction 做一次前后对账不保证保留窗口内全部中间值full-compaction两次大盘点之间做净变化对比不适合所以“完整 before/after”和“保留每一个中间事件”是两件事。十三、值没变要不要再发一组-U/U还有一种不太显眼的浪费同一个 Key 被重复写入但比较字段没有变化。默认情况下lookup和full-compaction的valueEqualiser可能为null源码条件是valueEqualisernull||!valueEqualiser.equals(before.value(),after.value())这意味着没有启用去重比较时即使 before 和 after 内容相同也可能生成一组-U/U。可以配置changelog-producer.row-deduplicatetruePaimon 会创建RecordEqualiser相同记录不再输出更新对。如果某些字段天然会变化但业务下游并不关心例如采集时间可以再配置忽略字段changelog-producer.row-deduplicate-ignore-fieldsingest_time这两个选项只对lookup和full-compaction有效。input的职责是保留上游原始输入不会在这里替上游去重。十四、ChangelogMode.all()不等于“每次 UPDATE 一定有四种记录”Flink 在构建执行计划前会询问 Source 能提供哪些 RowKind。BaseDataTableSource.getChangelogMode()的关键分支是if(!unbounded){returnChangelogMode.insertOnly();}if(table.primaryKeys().isEmpty()){returnChangelogMode.insertOnly();}if(mergeEngineFIRST_ROW){returnChangelogMode.insertOnly();}if(scanRemoveNormalize){returnChangelogMode.all();}if(options.get(CHANGELOG_PRODUCER)!NONE){returnChangelogMode.all();}returnChangelogMode.upsert();这里表达的是 Source 的能力边界场景向 Flink 声明的 ChangelogMode批读取Insert Only无主键追加表Insert OnlyFirst Row 主键表Insert Only主键表 noneUpsert主键表 input/lookup/full-compactionAllAll的意思是“这条 Source 链路允许出现四种 RowKind”不是承诺每一条 UPDATE 都一定同时出现-U/U。源码里还有一个scan.remove-normalizetrue的特殊分支它会让 Source 直接声明All。这是在调整 Flink Normalize 的边界不是 Producer 的常规选择本文都按默认值false来讨论。例如input模式下MySQL CDC Action 明明输入的是-D/IPaimon 就会输出-D/I。Source 声明All只是让 Flink Planner 知道这里可能出现撤回记录不能按纯追加流处理。十五、Source 是怎样找到 Changelog 文件的上一篇已经见过FollowUpScanner。到了这一篇那个分支就更好理解了switch(changelogProducer){caseNONE:followUpScannernewDeltaFollowUpScanner();break;caseINPUT:caseFULL_COMPACTION:caseLOOKUP:followUpScannernewChangelogFollowUpScanner();break;}ChangelogFollowUpScanner会先检查 Snapshot 是否带有snapshot.changelogManifestList()有 Changelog Manifest 才使用ScanMode.CHANGELOG没有就跳到下一个 Snapshot。这解释了一个常见现象full-compaction表明明不断有新 Snapshot下游却没有立刻收到变化。因为普通 Delta Snapshot 没有 Changelog ManifestSource 会继续等到产生 Changelog 的 Full Compaction Snapshot。十六、到底应该选哪个 Producer没有一个配置能同时做到“零额外成本、零延迟、完整旧值、保留每一次中间变化”。选择时先问下游真正需要什么。你的场景更合适的选择原因下游按主键覆盖只关心最终状态none写入成本最低Upsert 已经够用要保留 MySQL CDC 传来的逐笔变化input原样保存输入中间变化不容易被合并下游需要完整-U/U又希望较快看到lookup通过查旧值生成逻辑 before/after能接受批次延迟只关心稳定版本间净变化full-compaction在 Full Compaction 中统一比较下游做无主键聚合但不想让 Flink 维护大量 Normalize 状态lookup或full-compaction把旧值生成工作放到 Paimon 写入侧对于 MySQL CDC 入仓表如果你选择input一定要先接受一个事实你得到的是 MySQL CDC Action 的 -D/I不一定是 -U/U。如果下游协议强制要求 UPDATE_BEFORE/UPDATE_AFTER就应该考虑lookup、full-compaction或者确认 Flink Planner 会通过 Normalize 补齐。十七、上线以后重点盯这几个问题1.input写放大是否可以接受它会同时写数据文件和 Changelog 文件。CDC 更新频繁、宽表字段很多时Changelog 存储、Manifest 数量和文件清理压力都要计算进去。2.lookup是否拖慢 Checkpoint 或 Commitlookup-waittrue是默认值。热点 Bucket、随机主键 Lookup、慢存储都可能让 Compaction 成为写入链路的瓶颈。3.full-compaction延迟是否符合下游 SLA不要只看 Snapshot 是否增长要看 Snapshot 有没有changelog_manifest_list。可以查询SELECTsnapshot_id,commit_kind,commit_time,changelog_manifest_list,changelog_record_countFROMorders$snapshotsORDERBYsnapshot_idDESC;如果 Snapshot 一直有changelog_manifest_list却一直为空下游没有变化不一定是 Source 卡住也可能是 Producer 还没产出 Changelog。4.write-onlytrue是否配了独立 Compaction 作业write-only会让写入任务不负责 Compaction。对依赖 Compaction 生成日志的lookup和full-compaction来说如果没有独立 Compact ActionChangelog 就可能迟迟不出现。5. 不要期待改配置后补出旧历史Producer 影响的是后续写入和 Compaction。把none改成input不会倒回去给历史 Snapshot 补一套原始 Changelog。运行中的 Writer 和 Source 也不会因为表属性变了就自动重建内部链路。生产上调整 Producer应按一次作业变更处理评估历史边界、从 Savepoint 重启相关任务并验证新 Snapshot 是否开始产生 Changelog Manifest。十八、最后陪订单1001再走一遍作业第一次从S10做 latest-fullI[1001, 80.00, CREATED]MySQL 把金额改成100MySqlRecordParser先生成-D[1001, 80.00, CREATED] I[1001, 100.00, CREATED]如果是input原样写入 Changelog 文件 下游看到 -D / I如果是noneWrite Buffer 合并同一个 Key 增量数据文件只留下新的 Upsert 结果 下游不直接拿到旧值如果是lookup按 id1001 找到历史值 80 合并出当前值 100 比较 before / after 输出 -U[80]、U[100]如果是full-compaction先等待下一次 Full Compaction 比较最高 Level 中的 80 和本轮最终的 100 输出 -U[80]、U[100]到这里四个符号就不需要背了。只要顺着三个问题往下问即可这条记录是在加入结果还是撤回结果 旧值是上游传来的还是 Paimon 查出来的 我看到的是原始 Changelog还是 Flink Normalize 后的结果这三个问题想清楚大多数 Changelog “怎么和预期不一样”的问题就已经定位了一半。本篇关键源码位置RowKind.java定义I/-U/U/-D以及isAdd()、isRetract()CoreOptions.java定义changelog-producer四种模式和默认值BaseDataTableSource.java向 Flink 声明 Insert Only、Upsert 或 All ChangelogModeDataTableStreamScan.java根据 Producer 选择 Delta 或 Changelog FollowUpScannerDeltaFollowUpScanner.javanone模式读取 APPEND Snapshot 的 DELTAChangelogFollowUpScanner.java读取带 Changelog Manifest 的 SnapshotMergeTreeWriter.javainput模式创建 Changelog WriterSortBufferWriteBuffer.java把原始输入交给 Changelog Writer把合并结果交给数据文件 WriterDeduplicateMergeFunction.java主键表默认只保留同 Key 最新记录LookupChangelogMergeFunctionWrapper.javaLookup 旧值并生成I/-D/-U/UFullChangelogMergeFunctionWrapper.javaFull Compaction 中比较最高层旧状态与合并后新状态FullChangelogMergeTreeCompactRewriter.java只在 Full Compaction 输出到最高 Level 时生成 ChangelogMySqlRecordParser.javaMySQL before 转 DELETEafter 转 INSERTPrimaryKeySimpleTableTest.java验证 input 原样保留四种输入 RowKindPrimaryKeyFileStoreTableITCase.java验证 lookup 和 full-compaction 的完整 Changelog 输出本文源码基于 Apache Paimon 1.4.2。不同版本的 Changelog Producer、Compaction 和 Flink Planner 行为可能变化升级时请以对应版本源码和执行计划为准。
返回列表