Skip to content
P · Apache Paimon 源码第 5 课⏱ 15 分钟

05 提交:FileStoreCommitImpl 的冲突检测与重试

Apache Paimon 源码学习

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

原子抢快照、冲突检测与重试
原子抢快照、冲突检测与重试AI 生成配图

源码:operation/FileStoreCommitImpl.java、operation/commit/ConflictDetection.java、PrimaryKeyConflictDetection.java、CommitScanner.java、CommitRollback.java、catalog/RenamingSnapshotCommit.java、manifest/FileEntry.java

学习目标 ​

  • 理解乐观并发提交模型与一次提交的 4 个步骤。
  • 分清“抢快照号失败(重试)”和“真实冲突(抛异常)”。
  • 能解释任意一条提交冲突异常。

一句话结论 ​

不锁表,先写好 manifest,再尝试原子创建 snapshot-(N+1),谁先创建谁赢;输的一方重读最新快照、做冲突检测后重试。抢号失败会重试,文件冲突直接抛异常(作业 failover)。

1. 整体流程 ​

commit(ManifestCommittable)
 ├─ collectChanges() → 拆出 append 变更 / compact 变更
 ├─ APPEND  → tryCommit(kind=APPEND,  detectConflicts=checkAppendFiles)   ← 快照 1
 └─ COMPACT → tryCommit(kind=COMPACT, detectConflicts=true)               ← 快照 2

tryCommit()  重试循环
 while (true) {
   latest  = snapshotManager.latestSnapshot()
   changes = changesProvider.provide(latest)       // overwrite 会基于 latest 重算要删的文件
   result  = tryCommitOnce(retryResult, ..., latest)
   成功 → break
   超时 / 超次数 → 抛异常
   retryWaiter.retryWait(retryCount)               // 指数退避
 }

tryCommitOnce()  单次尝试
 ① 幂等检查(上次其实成功了吗?)
 ② 冲突检测
 ③ 构建 snapshot(合并 manifest,写 base/delta/changelog/index manifest)
 ④ 原子提交:commitSnapshotImpl → SnapshotCommit.commit

为什么一个 checkpoint 可能产生两个快照:append 只新增文件,新增不会和别人冲突,默认不检测;compact 要删除旧文件,必须检测。若 APPEND 中含 DELETE 条目或 DV 索引文件(如 Spark 行级 DELETE/UPDATE),shouldBeOverwriteCommit 会升级为 OVERWRITE,强制检测并允许回滚。

2. 重试循环 ​

抢号失败会退避重试,不等于冲突
抢号失败会退避重试,不等于冲突AI 生成配图
java
if (System.currentTimeMillis() - startMillis > options.commitTimeout()
        || retryCount >= options.commitMaxRetries()) {
    throw new RuntimeException("Commit failed ... there maybe exist commit conflicts between multiple jobs.");
}
retryWaiter.retryWait(retryCount);
配置默认
commit.max-retries10
commit.timeout无(Long.MAX_VALUE)
commit.min-retry-wait / commit.max-retry-wait10ms / 10s

退避:min(10ms × 2^n, 10s) + 最多 20% 随机抖动,避免多作业同时重撞。

会重试的结果只有:CommitFailRetryResult(抢号失败或原子提交异常)和 RetryCommitResult.forRollback(回滚了冲突的 compact 快照)。

3. tryCommitOnce 四步 ​

删除同一旧文件会触发真实冲突
删除同一旧文件会触发真实冲突AI 生成配图

① 幂等检查 ​

上一次原子提交抛异常(如网络超时)时,快照可能其实已写成功:

java
for (long i = startCheckSnapshot; i <= latestSnapshot.id(); i++) {
    Snapshot snapshot = snapshotManager.snapshot(i);
    hasOverwriteSinceLastAttempt |= snapshot.commitKind() == CommitKind.OVERWRITE;
    if (snapshot.commitUser().equals(commitUser)
            && snapshot.commitIdentifier() == identifier
            && snapshot.commitKind() == commitKind) {
        return new SuccessCommitResult();
    }
}

commitUser + commitIdentifier + commitKind 唯一标识一次提交(Flink 中 identifier 就是 checkpointId)。

② 冲突检测 ​

条件:latestSnapshot != null && (detectConflicts || discardDuplicate)。

读 base(scanBaseDataFiles):

  • 首次:readAllEntriesFromChangedPartitions,只读本次涉及的分区的全部文件清单。
  • 重试:复用上次的 baseDataFiles + 期间新快照的 delta(readIncrementalChanges)合并,重试很便宜;期间若发生 OVERWRITE 则退回全量读。
  • commit.discard-duplicate-files=true(仅 APPEND):从 delta 中剔除 base 里已存在的文件,用于幂等重复提交。

ConflictDetection.checkConflicts 四道检查:

a. checkBucketKeepSame          同分区 bucket 数必须一致(OVERWRITE 除外)
b. FileEntry.mergeEntries(delta)   delta 自身不能矛盾(同文件 ADD 两次等)
c. FileEntry.mergeEntries(base + delta) → checkDeleteInEntries
d. checkTableSpecificConflicts  主键表:LSM 层级重叠检查;append 表:无

核心 c 依赖合并规则:

java
case ADD:    checkState(old == null, "Trying to add file %s which is already in the map")
             map.put(id, entry);
case DELETE: if (map.containsKey(id)) map.remove(id);    // 先 ADD 后 DELETE,抵消
             else map.put(id, entry);                     // 孤立的 DELETE

合并后还残留 DELETE,说明“我要删的文件在最新快照里已经不存在”——被别人删了 → “File deletion conflicts detected! Give up committing.”(若残留全在过期分区,报 “You are writing data to expired partitions”)。

主键表 d:PrimaryKeyConflictDetection

java
按 (partition, bucket, level≥1) 分组,组内按 minKey 排序,相邻 a.maxKey >= b.minKey → “LSM conflicts detected!”

保护“同层文件不重叠”。场景:两个作业各自合并同一 bucket 的不同 L0 文件都输出到 L1,互不删对方文件(c 通过),但 L1 文件 key 范围重叠。

③ 构建 snapshot ​

  • 读上一快照所有 data manifest,ManifestFileMerger.merge,或复用上次尝试的合并结果(tryReuseManifestMergeResult)。
  • 写 baseManifestList、deltaManifestList、changelogManifestList、index manifest;继承 watermark、统计信息,构造 Snapshot。
  • 此阶段失败:肯定还没提交,commitCleaner 清理临时 manifest 后抛出。

④ 原子提交 ​

实现机制
RenamingSnapshotCommitfileIO.tryToWriteAtomic(snapshot-N);对象存储(rename 非原子)配合外部 Lock 并先检查 exists
CatalogSnapshotCommit由 Catalog(如 REST Server)基于 base 快照 uuid 做 CAS,见 11 章

三种结果:

结果含义处理
true抢到 snapshot-N成功,触发 commit callbacks(Hive 分区同步、自动 tag、Iceberg 等)
false被别人抢了 NforCommitFail(...),保留临时文件,重试时复用 base 与 manifest 合并结果
异常不确定是否成功不清理,返回 forCommitFail(..., e, ...),下次由①判定

4. 回滚特例:DML vs 后台合并 ​

java
if (exception.isPresent()) {
    if (allowRollback && rollback != null && rollback.tryToRollback(latestSnapshot))
        return RetryCommitResult.forRollback(exception.get());
    throw exception.get();
}

CommitRollback.tryToRollback:只有最新快照是 COMPACT 时才回滚到 latest-1。场景:Spark DELETE 与流作业后台合并重写了相同文件,放弃合并结果(合并可重做),让用户 DML 重试成功。

5. 重启后的幂等:filterCommitted ​

java
latestSnapshot = snapshotManager.latestSnapshotOfUser(commitUser);
只保留 committable.identifier() > latestSnapshot.commitIdentifier() 的

与①配合实现 exactly-once。

6. 常见场景 ​

场景结果
两个作业只 APPEND 同一张表只会抢号,自动重试成功
两个作业写同一 bucket 且都做合并File deletion conflict,failover。建议写入作业 write-only=true + 独立 compaction 作业
两作业各自合并不同 L0 到 L1LSM conflict
从旧 savepoint 恢复要删的文件早被合并掉 → 持续冲突。从最新 savepoint 恢复,或把表回滚到对应快照
不 overwrite 就改 bucket 数checkBucketKeepSame 失败,需要 INSERT OVERWRITE 重分布
Spark DELETE vs 流合并回滚 compact 快照,DML 重试成功

冲突异常中会打印 “Base commit user / Current commit user”:不同说明有别的作业在写同一分区;相同通常是从旧 checkpoint 恢复。

动手实验 ​

  • [ ] 跑 FileStoreCommitTest 中搜 conflict 的用例。
  • [ ] 自己写一个测试:两个 commitUser 对同一 bucket 的同一批文件做 compact 提交,复现 File deletion conflict。
  • [ ] 日志搜 Atomic commit failed for snapshot(抢号失败)和 conflicts detected(真实冲突)。

自测题 ​

  1. 为什么 APPEND 默认不做冲突检测?什么情况下会被升级为 OVERWRITE?
  2. mergeEntries 后残留 DELETE 意味着什么?
  3. 为什么还需要 LSM 层级重叠检查?举一个 c 检查发现不了的例子。
  4. 原子提交抛异常时为什么不能清理临时文件?下次如何判断成功与否?
  5. 回滚为什么只针对 COMPACT 快照?

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