Skip to content

S2 第 8 讲:Lookup Changelog——-U/+U 是怎么算出来的 ​

Apache Paimon 源码学习

作者 X老师(DaemonforY),Paimon 2.0 / master 源码,按 CC BY-NC-SA 4.0 发布。配套实验和代码在 GitHub。

Yui 和 Kai 在主键工坊追踪旧值与新值。
Yui 和 Kai 在主键工坊追踪旧值与新值。AI 生成配图

TL;DR:第一季第五讲(下)讲过,下游按某个值做聚合时,需要 -U 告诉它旧值是什么,才能把旧值撤回。Paimon 主键表默认(changelog-producer = none)不存这种变更,流读时 Flink 要加一个 ChangelogNormalize 算子,用状态记住每个 key 的最新值,自己补出 -U。lookup 把这件事挪到了写入时:每次提交前,把新写的 L0 文件强制合并到更高层;合并时对每个 key 去高层查(lookup)旧值当 before,合并结果当 after,比较二者写出 changelog 文件。实验里第二次写入 2 行,产出了 4 条 changelog——值没变的那一行也发了一对 -U/+U,除非打开 changelog-producer.row-deduplicate。

实验:四张表,同样三步 ​

labs/sql/lab10/lookup-changelog.sql 建四张 bucket = 1 的主键表,做同样的写入:

① INSERT (1,'a') (2,'b') (3,'c')
② INSERT (1,'A') (2,'b')        -- key 1 改了值,key 2 写入的值和原来一样
③ DELETE FROM … WHERE k = 3
表选项
cl_none默认(changelog-producer = none)
cl_lookupchangelog-producer = lookup
cl_deduplookup + changelog-producer.row-deduplicate = true(只做 ① ②)
cl_agglookup + merge-engine = aggregation,total 求和(写 (1,10) (2,20),再写 (1,5))

每一步之后看 $snapshots、$files,再用 $audit_log 加 incremental-between 读出某两个快照之间的变更(rowkind 一列就是 +I/-U/+U/-D)。

四张表读到的变更

默认表 cl_none 读第 ② ③ 步的增量,拿到的是写进去的原始记录(节选自 run-2.0.0-flink2.2.log):

Flink SQL> SELECT * FROM `cl_none$audit_log` /*+ OPTIONS('incremental-between' = '1,3') */;
|                             +I |           1 |                              A |
|                             +I |           2 |                              b |
|                             -D |           3 |                              c |
3 rows in set

没有 -U:下游不知道 key 1 原来是 a。cl_lookup 第 ② 步只写了 2 行,读出来的是:

Flink SQL> SELECT * FROM `cl_lookup$audit_log` /*+ OPTIONS('incremental-between' = '2,4', 'incremental-between-scan-mode' = 'changelog') */;
|                             -U |           1 |                              a |
|                             +U |           1 |                              A |
|                             -U |           2 |                              b |
|                             +U |           2 |                              b |
4 rows in set

a 从哪来?第 ② 步写进去的数据里根本没有 a。下面分三个问题拆:什么时候算、每个 key 怎么算、算完写到哪。

一、什么时候算:每次提交前,把 L0 强制推上去 ​

提交前先压实低层文件,再产生变更记录。
提交前先压实低层文件,再产生变更记录。AI 生成配图

什么时候算

先看 cl_lookup$snapshots(节选自 run-2.0.0-flink2.2.log,最后一次查询):

|          snapshot_id |                    commit_kind |   total_record_count |   delta_record_count | changelog_record_count |
|                    1 |                         APPEND |                    3 |                    3 |                 <NULL> |
|                    2 |                        COMPACT |                    3 |                    0 |                      3 |
|                    3 |                         APPEND |                    5 |                    2 |                 <NULL> |
|                    4 |                        COMPACT |                    5 |                    0 |                      4 |
|                    5 |                         APPEND |                    6 |                    1 |                 <NULL> |
|                    6 |                        COMPACT |                    6 |                    0 |                      1 |

每次写入都变成两个快照:一个 APPEND(新写的 L0 文件),一个 COMPACT,changelog 只挂在 COMPACT 快照上(3、4、1 条)。对照 cl_none:三次写入只有 3 个 APPEND 快照,changelog_record_count 都是 NULL,三个文件都还在 L0——3 个 sorted run 没到 UniversalCompaction 的触发条件(第 3 讲,默认 5 个),根本不合并。

cl_lookup 为什么每次都合并?开启 lookup 后,合并策略换成了 ForceUpLevel0Compaction(MergeTreeCompactManagerFactory.java 第 220~236 行):先按 Universal 的规则挑,挑不出来就 forcePickL0,把 L0 强制合并上去(ForceUpLevel0Compaction.pick,第 50~77 行)。jdb 在三次写入时各命中一次 forcePickL0,当时的 sorted run 数分别是 1、2、3。

三次合并之后的文件(cl_lookup$files):

步骤文件说明
① 后L5 · 3 行只有一个 run,直接推到最高层
② 后L4 · 2 行、L5 · 3 行只合并了 L0 这一个文件;输出层 = 下一个更老的 run(L5)减 1(第 3 讲)
③ 后L3 · 1 行、L4 · 2 行、L5 · 3 行L3 里是 key 3 的删除记录:没到最高层,不能丢

注意第 ② 步:这次合并的输入只有新写的 L0 文件,L5 并不参与。那 key 1 的旧值 a 只能到 L5 里去查——这就是 “lookup” 的意思。

流作业里还有一个开关 lookup-wait(默认 true,CoreOptions 第 2284~2289 行):Flink sink 把“需要 lookup 且 lookup-wait”作为每次 prepareCommit 的 waitCompaction(StoreSinkWrite.java:146、CoreOptions.prepareCommitWaitCompaction 第 4616~4622 行),checkpoint 会等这次合并做完,数据文件和 changelog 文件进同一次提交。本讲的实验是批作业,批作业结束时本来就会等合并(PrepareCommitOperator.java:101-103),所以实验本身证明不了 lookup-wait 的作用,这一段是读源码得出的。

二、每个 key 怎么算:before 去高层查,after 是合并结果 ​

旧值去高层查,新值由合并结果得到。
旧值去高层查,新值由合并结果得到。AI 生成配图

每个 key 怎么算

合并时同一个 key 的记录交给 LookupChangelogMergeFunctionWrapper,核心是 getResult()(第 104~146 行),节选:

java
// 1. 本次参与合并的记录里,挑出 level > 0 中最新的一条
KeyValue highLevel = mergeFunction.pickHighLevel();
boolean containLevel0 = mergeFunction.containLevel0();
// 2. 没有,就去更高层查
if (highLevel == null) {
    T lookupResult = lookup.apply(mergeFunction.key());   // = lookupLevels.lookup(key, outputLevel + 1)
    ...
}
// 3. 合并出新值
KeyValue result = mergeFunction.getResult();
// 4. 有 L0 记录参与时才产 changelog
if (containLevel0 && lookupStrategy.produceChangelog) {
    setChangelog(highLevel, result);
}

lookup 从 outputLevel + 1 开始(LookupMergeTreeCompactRewriter.java:222),逐层往下找,第一个命中就返回(LookupUtils.lookup,第 35~57 行)。为什么不用看更老的层?因为离开 L0 的数据都已经和当时查到的旧值合并过,L1 及以上每个 key 的记录都是完整的最新状态,最上面那条就是 before。

setChangelog(before, after)(第 148~163 行)的规则:

beforeafterchangelog
没有 / 是删除新增+I
有删除-D
有新增,值不同-U + +U
有新增,值相同默认仍是 -U + +U;开 row-deduplicate 后不输出

jdb 在 getResult 的第 137 行(after 算出来之前)和 setChangelog 入口各下一个断点,cl_lookup 三步的打印(节选自 jdb-lookup.log,KeyValue@… 是对象地址,表示 before 存在):

#3   getResult   key = 1   highLevel = null              containLevel0 = true
#4   setChangelog          before = null                 after = INSERT "a"      → +I
#11  LookupLevels.createLookupFile   data-ec659df9…  level = 5
#12  getResult   key = 1   highLevel = KeyValue@52b55058 containLevel0 = true
#13  setChangelog          before = KeyValue@52b55058    after = INSERT "A"      → -U / +U
#14  getResult   key = 2   highLevel = KeyValue@4ee3f48c containLevel0 = true
#15  setChangelog          before = KeyValue@4ee3f48c    after = INSERT "b"      → -U / +U
#19  getResult   key = 3   highLevel = KeyValue@654ca2a1 containLevel0 = true
#20  setChangelog          before = KeyValue@654ca2a1    after = DELETE "c"      → -D

(为便于阅读,整理成每行一次命中,省略线程名;→ 后面是按上表规则对应的 changelog,和 $audit_log 读出的结果一致。)

第 ② 步的第 11 次命中很说明问题:在拿到 key 1 的 before 之前,先为 L5 的数据文件在本地建了一个 lookup 文件(LookupLevels.java:206)。lookup 不是每次去扫 Parquet,而是把高层数据文件转换成本地的 KV 文件(createLookupFile,第 308 行起)再点查,并缓存起来;第 ③ 步(第 18 次命中)又为同一个 L5 文件建了一次——这个实验里每条 INSERT/DELETE 都是一个独立的批作业,缓存跟着作业走。流作业里同一个写入算子会复用它(读源码得出,未单独实验)。

反直觉的一点:key 2 第 ② 步写入的 b 和原来的 b 一模一样,也发了一对 -U b +U b。只有 changelog-producer.row-deduplicate = true(默认 false,CoreOptions 第 1132~1136 行)时,才会创建 valueEqualiser 比较前后值(KeyValueFileStore.java:85-89)。cl_dedup 第 ② 步读出来只剩两条:

Flink SQL> SELECT * FROM `cl_dedup$audit_log` /*+ OPTIONS('incremental-between' = '2,4', 'incremental-between-scan-mode' = 'changelog') */;
|                             -U |           1 |                              a |
|                             +U |           1 |                              A |
2 rows in set

上游经常重复写同样的数据(比如 CDC 全量 + 增量重放)时,这个开关能省下大量无意义的变更。

三、算完写到哪:changelog 文件 + 三种升级策略 ​

三种升级策略

changelog 总是写成单独的 changelog-* 文件(ChangelogMergeTreeRewriter.rewriteOrProduceChangelog,第 125~210 行),cl_lookup 的 bucket 目录下三步后有 3 个 changelog 文件。不同的是数据文件要不要重写。L0 文件“升级”(只有它自己、不和别的文件重叠)时,LookupMergeTreeCompactRewriter.upgradeStrategy()(第 139~173 行)有三种结果,jdb 全部命中过:

情况策略jdb实验证据
输出到最高层CHANGELOG_NO_REWRITE(第 160 行)各表第 ① 步,outputLevel = 5没有更老的数据,文件内容就是完整状态
去重引擎、没有 sequence.fieldCHANGELOG_NO_REWRITE(第 166 行)cl_lookup ② ③,outputLevel = 4、3第 ② 步的 L0 文件 bb385d56 升到 L4,文件名不变
其它引擎(aggregation 等)CHANGELOG_WITH_REWRITE(第 172 行)cl_agg ②,outputLevel = 4L0 文件 84d107bc → L4 文件 94f15597,换了新文件

“NO_REWRITE” 的意思是:读一遍这个文件算出 changelog,但数据文件不重写,只把它的 level 改掉。去重引擎可以这样做,因为新值就是最终结果;聚合引擎不行——cl_agg 第 ② 步的 L0 里只有增量 (1, 5),如果原样升到 L4,L4 存的就是 5,破坏了“L1+ 是完整状态”的不变式。所以它必须把合并后的 (1, 15) 写成新文件,查询结果是 15,changelog 是 -U 1 10 +U 1 15。

(文件名来自 run-2.0.0-flink2.2.log 中按快照 3、快照 4 查询的 $files;jdb 那次运行的文件名不同,见 jdb-lookup.log 第 2、10、17、43 次命中。)

四、和 none / input / full-compaction 怎么选 ​

四种 changelog-producer

changelog-producer-U 从哪来代价
none(默认,CoreOptions 第 1105~1111 行)流读时 Flink 加 ChangelogNormalize,用状态记住每个 key 的最新值写入最省;下游状态随 key 数增长
input直接把输入当 changelog,要求源本身是完整的 CDC 流低;输入不完整就会错
lookup本讲每次提交都合并;本地 lookup 文件;checkpoint 等合并
full-compaction只在全量合并时对比新旧(FullChangelogMergeFunctionWrapper)写入较省;changelog 延迟取决于全量合并间隔

在 Flink 2.2.0 上流读同样的数据(实验 3,run-all 日志),两种表的执行计划差一个算子:

none 表:  ChangelogNormalize(key=[order_id])
           +- Exchange(distribution=[hash[order_id]])
              +- TableSourceScan(table=[[paimon, default, orders_none]], …)
lookup 表:没有 ChangelogNormalize,直接 TableSourceScan(table=[[paimon, default, orders_lookup]], …)

(节选自 lab03-a-read-none-flink2.2.log / lab03-a-read-lookup-flink2.2.log。)这就是 lookup 的取舍:用写入时的合并和 lookup,换掉下游每个流作业里的 ChangelogNormalize 状态。读的作业越多、key 越多,越划算。还没做性能压测,具体代价要在自己的数据量上测。

小结

  1. 时机:L0 上推的那一刻。lookup 模式下合并策略是 ForceUpLevel0Compaction,每次提交前都把 L0 推上去,流作业的 checkpoint 会等它(lookup-wait)。
  2. 算法:before = 高层 lookup 到的旧值(不变式保证它是完整状态),after = 合并结果,按表格出 +I / -U+U / -D。
  3. 输出:changelog 写成单独文件、挂在 COMPACT 快照上;去重引擎升级只改 level,其它引擎要重写。
  4. 值没变也发 -U/+U,需要时开 row-deduplicate。

下一讲:Flink 写入的两阶段提交——这些 changelog 文件和数据文件,是怎么借 checkpoint 做到 exactly-once 地提交进同一个快照的?(回扣第一季第八讲下篇 Sink V2。)

自测题:cl_lookup 如果第 ④ 步再写一次 (1, 'A'),会读到什么 changelog?这次 lookup 会在哪一层命中?(答案:before 是 L4 里的 (1, A)——lookup 从 outputLevel + 1 开始逐层找,L4 在 L5 之前;值相同,默认仍发 -U 1 A +U 1 A。具体输出层取决于当时的 sorted run,未单独运行,可用同一断点清单验证。)


代码示例在页面里运行时使用 HiveGPT 的模型接口。延伸阅读来自 JavaGuide(Apache-2.0),版权归原作者。