Skip to content

S2 第 11 讲:Spark MERGE INTO 的三条路径 ​

Apache Paimon 源码学习

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

Yui和Kai见证数据合并后分流到三条写入路径
Yui和Kai见证数据合并后分流到三条写入路径AI 生成配图

TL;DR:Paimon 的 Spark MERGE INTO 有三条落地方式,但前半段是同一个引擎:source 和 target 做 full outer join,每一行按 MATCHED / NOT MATCHED / NOT MATCHED BY SOURCE 归类,找第一个成立的动作,输出 +I / +U / -D。区别在后半段:主键表只写变化的行(新 L0 文件),旧版本交给 LSM,提交是 APPEND;追加表把被触及的文件整个重写(copy-on-write),提交升级为 OVERWRITE;追加表 + 删除向量旧文件一个字节不动,按行号打标记,也是 OVERWRITE。默认三者都走 Paimon 自己的 V1 命令 MergeIntoPaimonTable;只有打开 spark.paimon.write.use-v2-write,追加表才交给 Spark 原生的 V2 copy-on-write。

实验:同一条 MERGE,四张表 ​

这一讲第一次用 Spark:labs 新增了一个 spark/ 子工程(Spark 3.5.8 + paimon-spark-3.5),./spark.sh <sql 文件> 用法和 Flink 的 run.sh 一样,在本地 local[1] 跑(单线程,文件数每次相同)。

labs/sql/lab13/merge-into.sql:每张表先分两批写 (1,a)(2,b)(3,c) 和 (4,d)(5,e)(6,f),得到两个数据文件;然后执行同一条 MERGE:

sql
CREATE TEMPORARY VIEW s AS SELECT * FROM VALUES (2, 'B', 'U'), (3, 'C', 'D'), (7, 'G', 'I') AS s(id, v, op);

MERGE INTO t USING s ON t.id = s.id
  WHEN MATCHED AND s.op = 'D' THEN DELETE
  WHEN MATCHED THEN UPDATE SET v = s.v
  WHEN NOT MATCHED THEN INSERT (id, v) VALUES (s.id, s.v);

即:id 2 改成 B、删除 id 3、插入 id 7。四张表:主键表 t_pk、追加表 t_cow、追加表 + 删除向量 t_dv,以及打开 V2 写的追加表 t_v2。四张表 MERGE 后的查询结果完全一样:1 a、2 B、4 d、5 e、6 f、7 G。

三张表的文件变化

每行在哪个文件、第几行,用 Paimon 的元数据列 __paimon_file_path、__paimon_row_index 直接查(节选,文件名只取前 8 位):

t_cow,MERGE 之前                       t_cow,MERGE 之后
|1  |a  |28525a1c|                      |1  |a  |cc9c094d|
|2  |b  |28525a1c|                      |2  |B  |cc9c094d|
|3  |c  |28525a1c|                      |4  |d  |8209d23b|
|4  |d  |8209d23b|                      |5  |e  |8209d23b|
|5  |e  |8209d23b|                      |6  |f  |8209d23b|
|6  |f  |8209d23b|                      |7  |G  |cc9c094d|
表MERGE 后的文件提交
主键表 t_pk原来两个文件不动,多一个 3 条记录的 L0 文件 a364b6ebAPPEND(total 9,delta 3)
追加表 t_cow含 1~3 的 28525a1c 被换成 cc9c094d(1 a、2 B、7 G),含 4~6 的 8209d23b 不动OVERWRITE(total 6,delta 0)
追加 + DV t_dv两个旧文件都不动;新增 adbc87b5(2 B、7 G),外加一个 33 字节的删除向量索引OVERWRITE(total 8,delta 2)

t_dv MERGE 之后:id 1 还在 2cd6c04f 的第 0 行,原来的第 1、2 行(id 2、3)查不到了——它们被记进了删除向量,文件本身没变。$snapshots 的 total 是 8(6 + 2),因为物理上 8 行都还在。

一、共用的引擎:full outer join + 逐行判定 ​

两股数据经全外连接后逐行分流判定
两股数据经全外连接后逐行分流判定AI 生成配图

共享引擎

三条 V1 路径都在 MergeIntoPaimonTable(paimon-spark-common)里。run()(第 86~97 行):

scala
checkMatchRationality(sparkSession)            // 一行 target 不能匹配多行 source
...
val commitMessages = if (withPrimaryKeys) {
  performMergeForPkTable(sparkSession)
} else {
  performMergeForNonPkTable(sparkSession)
}
writer.commit(commitMessages, Snapshot.Operation.MERGE)

核心是 constructChangedRows:target 加一列 _target_row_ = true,source 加一列 _source_row_ = true,做 full outer join,然后逐行看:

join 结果归类本例输出
两边都有MATCHEDid 3(op = 'D')→ DELETE;id 2 → UPDATE主键表 / DV:-D 3、+U 2 B;CoW:3 不写回、写 2 B
只有 sourceNOT MATCHEDid 7 → INSERT+I 7 G
只有 targetNOT MATCHED BY SOURCEid 1、4、5、6(本例没有这类动作)来自被重写的文件(_file_touched_col_ = true)就原样写回,否则丢弃

动作按 SQL 里写的顺序匹配,第一个成立的生效——所以 WHEN MATCHED AND s.op = 'D' THEN DELETE 要写在 WHEN MATCHED THEN UPDATE 前面。

先做一次检查(checkMatchRationality,第 341~360 行):只要有 MATCHED 动作,就先 inner join 统计每行 target 匹配了几行 source,有大于 1 的直接报错。实验里拿两行 id = 2 的 source 去 MERGE:

[EXPECTED ERROR] RuntimeException: Can't execute this MergeInto when there are some target rows that each of them match more than one source rows. It may lead to an unexpected result.

同一行被改成 X 还是 Y 取决于执行顺序,结果不确定,所以干脆拒绝。用 MERGE 前先对 source 按 key 去重。

二、三条路径怎么落地 ​

主键表、重写表和原生行级操作形成三条路径
主键表、重写表和原生行级操作形成三条路径AI 生成配图

主键表:performMergeForPkTable(第 100~106 行)把 constructChangedRows 的结果(remainDeletedRow = true,-D 也保留)直接交给 writer,按主键分 bucket 写成新的 L0 文件。不读也不重写任何旧文件——旧版本在读的时候按 sequence 合并掉(第 6 讲),合并时清除(第 4 讲);如果是 lookup 表,还会为流读算出 -U/+U(第 8 讲)。实验里新文件 a364b6eb 有 3 条记录;t_pk$audit_log 在快照 3 上看到 +U 2 B、-D 3 c、+I 7 G。(新文件里具体是哪 3 条,由 record_count 和 audit_log 推出,未逐行读出该文件。)

追加表(copy-on-write):performMergeForNonPkTable(第 108 行起),没开删除向量时:

  1. 按 target 上能下推的条件挑候选文件;
  2. 用 inner join 找出被触及的文件(有 MATCHED 的 UPDATE / DELETE)——本例是 28525a1c;只有 INSERT 时会找“要读但不用重写”的文件(第 152~191 行);
  3. 被触及的文件全部读出来,打上 _file_touched_col_ = true,走同一个引擎:没有动作命中的行(id 1)原样写回,删除的行(id 3)不写回;
  4. 新写一个文件,旧文件记为删除。

所以 id 1 一个字没改,也被拷贝进了新文件。文件越大、要改的行越分散,拷贝越多——这就是 CoW 的写放大。

追加表 + 删除向量:同一个函数里的 DV 分支(第 119~145 行):扫描时带上 __paimon_file_path、__paimon_row_index,引擎输出里 -D 和 +U 的行按“文件 + 行号”收集成删除向量(+U 的位置是被更新那行原来的位置);+I 和 +U 的行写成新文件。旧文件一个字节不动。

三、走哪条路:默认都是 V1 ​

路径选择

EXPLAIN MERGE INTO 的物理计划,三张表默认都是(节选):

== Physical Plan ==
Execute MergeIntoPaimonTable
   +- MergeIntoPaimonTable PrimaryKeyFileStoreTable[default.t_pk], …
Execute MergeIntoPaimonTable
   +- MergeIntoPaimonTable AppendOnlyFileStoreTable[default.t_cow], …
Execute MergeIntoPaimonTable
   +- MergeIntoPaimonTable AppendOnlyFileStoreTable[default.t_dv], …

追加表什么时候交给 Spark 原生的 V2 行级操作?SparkTable.supportsV2RowLevelOps(第 87~130 行)要求 Spark ≥ 3.5、useV2Write 为 true,并且没有主键、没有删除向量、没有 data evolution、没有 CHAR 列。useV2Write 来自 Spark 会话配置 spark.paimon.write.use-v2-write,默认 false(SparkConnectorOptions.java:91-96、OptionUtils.scala:55-68)。实验里打开它再对 t_v2 做同样的 MERGE:

SET `spark.paimon.write.use-v2-write`=true;

== Physical Plan ==
ReplaceData PaimonWrite(table=default.t_v2, overwriteDynamic=true)
+- AdaptiveSparkPlan isFinalPlan=false
   +- Project [id#1743, v#1744]
      +- MergeRowsExec[id#1743, v#1744, __paimon_file_path#1745]
         +- SortMergeJoin [id#1724], [id#1726], FullOuter
            …
            +- BatchScan default.t_v2[id#1724, v#1725, __paimon_file_path#1732] PaimonCopyOnWriteScan: [t_v2] RuntimeFilters: dynamicpruningexpression(__paimon_file_path#1732 IN subquery#175…

(节选,省略了子查询的展开。)换成了 Spark 的 ReplaceData:同样是 full outer join,但运行时过滤——先用一个子查询(target 与 source 的 LeftSemi join)找出含匹配行的文件,再只扫这些文件、整体替换。结果和 V1 CoW 一样:含 1~3 的 85c3af2f 被换成 615136a3(1 a、2 B、7 G),e046daef 不动,提交也是 OVERWRITE。

注意 SET 的写法:键里有 -,要用反引号包起来,否则 Spark 报 INVALID_SET_SYNTAX(实验第一次就踩了这个坑)。

四、提交:APPEND 还是 OVERWRITE ​

![提交类型

Paimon 提交日志(节选):

t_pk :Finished collecting changes, including: 1 append table files                     → kind APPEND
t_cow:Finished collecting changes, including: 2 append table files                     → kind OVERWRITE
t_dv :Finished collecting changes, including: 1 append table files, 1 append index files → kind OVERWRITE

t_cow 的“2 append table files”是 1 个新增、1 个删除(被重写的旧文件)。FileStoreCommitImpl 第 351~356 行:

java
if (conflictDetection.shouldBeOverwriteCommit(appendSimpleEntries, changes.appendIndexFiles)) {
    commitKind = CommitKind.OVERWRITE;
    checkAppendFiles = true;
    allowRollback = true;
}

shouldBeOverwriteCommit(ConflictDetection.java 第 191~204 行):这次提交的文件里有 DELETE,或者带了删除向量索引,就返回 true。于是:

  • 主键表只加新文件:APPEND,默认不做冲突检测,可以和 Flink 流写同时写一张表(第 5 讲);
  • 追加表要删旧文件或改删除向量:升级为 OVERWRITE,强制冲突检测;如果期间别的作业合并了同一批文件、最新快照是 COMPACT,允许回滚那个 COMPACT 快照再重试(allowRollback,教程 10 章第 8 节;本讲未复现并发场景)。

五、怎么选 ​

写放大与选择

本例写了几行旧数据适合
主键表3(只写变化)读时 / 合并时按主键消除频繁更新、和 Flink 流写共存
追加表 CoW3(其中 1 行原样拷贝)整个文件重写、旧文件删除偶尔批量修正、读多写少
追加表 + DV2 + 一个删除向量旧文件不动,按行号标记没有主键但经常改少量行

本例每个文件只有 3 行,CoW 只多拷贝 1 行;真实表一个文件几十万行时,改一行也要重写整个文件。没有做大表压测,这里只说明机制。

小结

  1. 三条路径共用 full outer join + 逐行判定,一行 target 匹配多行 source 直接报错。
  2. 主键表只写变化,APPEND 提交;追加表 CoW 重写被触及的文件,DV 只打标记,两者都升级为 OVERWRITE 并强制冲突检测。
  3. 默认都走 Paimon 的 V1 命令;spark.paimon.write.use-v2-write = true 时,满足条件的追加表交给 Spark 原生 V2,运行时只扫含匹配行的文件。

下一讲:REST Catalog——提交时谁来保证 snapshot-(N+1) 只有一个人写成功?第 5 讲说对象存储上 rename 不可靠,REST Catalog 是服务端的答案。

自测题:给 t_cow 的 MERGE 加一句 WHEN NOT MATCHED BY SOURCE THEN DELETE,会重写几个文件?(答案:两个文件都要重写。id 1、4、5、6 都只在 target 里,会被删除;含 1~3 的文件本来就因 MATCHED 被重写,含 4~6 的 8209d23b 现在也被触及,重写后没有剩下的行。最后表里只剩 2 B、7 G。而且有 NOT MATCHED BY SOURCE 时不能只靠 target 条件裁剪候选文件(教程 10 章第 10 节)。按源码推导,未单独运行。)


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