TT Lab
开始
学习 学习路径 课程

湖仓表格式 — 从元数据理解 Apache Iceberg

两个写入者同时提交同一张表 — 谁赢、留下什么、丢了什么

在 TT Lab 中继续学习

目标

用 pyiceberg 确定性地复现两个写入方读取同一个 metadata 之后依次提交的竞争。追加(append)通过自动重试,两边都能写入;关掉重试后,一方会失败,写了一半的文件会作为孤儿留下。带条件的覆盖写入(overwrite)即使重试也会被验证拦下,之后必须重新读取、重新计算,才不会出现丢失更新,这一点用一个计数器来确认。

为什么重要

Iceberg 没有锁。写入方各自写完文件、创建新的 metadata,然后在 Catalog 中尝试“如果我读到的 main 仍然没变,就改掉它”这个条件式替换。两个同时到达,只有一个获胜,输的一方重新读取新的 metadata,把自己的变更重新应用上去。这就是乐观并发控制——相信大多数情况下不会冲突,冲突了就重来。 问题在于,有“可以重新应用的变更”,也有“不可以这样做的变更”。添加新文件,无论这期间有人做了什么,都可以重新应用。但是“把符合这个条件的行改成这个值”这样的变更,如果这期间有人修改了同样条件的行,结果就会出错。库会在文件层面验证这一点并拒绝,但如果你的代码不用以前读到的值重新计算,而是原样重写,就没有任何人能拦住。

步骤

  1. 在 /root/ice/conc/common.py 中放一个用 pyarrow 读取一天 CSV 的辅助函数(order_ts 使用 UTC 时区),并用 /root/ice/conc/setup.py 创建 lake.conc.orders,写入 2026-03-01。
  2. 用 /root/ice/conc/race.py 先把两个 Table 对象 a 和 b 都读取出来,再让 a 追加 03-02,b 追加 03-03,并把结果写入 /root/ice/conc/out/race.json。
  3. 用 /root/ice/conc/noretry.py 把表属性 commit.retry.num-retries 设为 0,再次进行同样的竞争(03-04、03-05),然后把输的一方的异常名称写入 /root/ice/conc/out/noretry.txt。
  4. 用 /root/ice/conc/orphans.py 找出位于表位置的 data 目录中、但没有任何快照指向的 Parquet 文件,写入 /root/ice/conc/out/orphans.json。
  5. 用 /root/ice/conc/conflict.py 创建计数器表 lake.conc.counters(hits = 10),让两个工作者读取相同的值后,各自覆盖写入 +5,并把输的一方的异常名称写入 /root/ice/conc/out/conflict.txt。
  6. 用 /root/ice/conc/retry.py 正确地重做输的一方的 +5,让计数器变成 20。
  7. 在 /root/ice/conc/report.md 中写出 ## 자동 재시도、## 실패한 커밋의 흔적、## 잃어버린 갱신 三个小节。

参考

用 Python 创建的表

在 /root/ice/conc/common.py 中,让 day("YYYY-MM-DD") 把当天的 CSV 返回为 pyarrow 表(amount 为 int32,order_ts 为 UTC 时区的 timestamp),并用 /root/ice/conc/setup.py 以 format-version 2 创建 lake.conc.orders,写入 2026-03-01。

pyiceberg 的 create_table(이름, schema=pyarrow_스키마) 会为每一列分配字段 ID(占位符依次为表名与 pyarrow schema)。给时间加上时区,就会成为与 Spark 创建的表相同的 timestamptz。评分器会检查第一个快照是否加上了 3 月 1 日的行数。

竞争——追加重新应用就行

在 /root/ice/conc/race.py 中,先把a = load_table(…) 和 b = load_table(…) 都创建出来,再按 a.append(03-02)、b.append(03-03) 的顺序提交,并把快照数和行数按 {"snapshots", "rows"} 的格式写入 /root/ice/conc/out/race.json。

b 持有的是 a 提交之前的 metadata,所以第一次提交尝试会被条件拦住。追加无论这期间进来了什么,都可以在其上重新应用,所以 pyiceberg 会重新读取并再次尝试,最终成功。评分器会检查三个快照是否连成了一条线(没有分叉)以及行数。

关掉重试,输的一方就会失败

用 /root/ice/conc/noretry.py 把 lake.conc.orders 的属性 commit.retry.num-retries 改为 "0",然后像第 2 步那样先创建 a 和 b,让 a 追加 03-04,b 追加 03-05。把 b 的异常名称写入 /root/ice/conc/out/noretry.txt 的第一行。

重试为 0 时,一旦在条件式替换中落败,异常就会抛出来。然而 b 在尝试提交之前,已经把数据文件全都写好了。评分器会检查异常名称、属性值,以及表中是否有 3 月 4 日而没有 3 月 5 日。

失败的提交留下的文件

用 /root/ice/conc/orphans.py 把所有快照所指向的文件列表,与表位置(t.location())下 data 目录中实际的 Parquet 文件对比,把列表中没有的文件(孤儿)的路径,按 {"orphans": [경로, …]} 的格式写入 /root/ice/conc/out/orphans.json(占位符为文件路径)。

孤儿文件不是表的一部分,所以读不到,但会占用空间。必须与所有快照所指向的文件对比,而不是只与当前快照对比——旧快照的文件,是为时间旅行而存活的文件。清理这类文件的,就是下一个模块的 remove_orphan_files。

覆盖写入会被验证拦下

用 /root/ice/conc/conflict.py 以 hits = 10 这一行新建(如已存在则先删除)lake.conc.counters(name STRING, value BIGINT),让 a 和 b 都读取值之后,a 用 overwrite(…, overwrite_filter=EqualTo("name", "hits")) 写入 읽은 값 + 5(占位符为读到的值),b 也以同样的方式写入。把 b 的异常名称写入 /root/ice/conc/out/conflict.txt 的第一行。

b 的第一次尝试在条件式替换中落败,重试时,会被“这期间有符合我的条件(name = hits)的文件新进来了”这项验证拦下而中止。因为与追加不同,覆盖写入如果原样重新应用,就会抹掉 a 的结果。评分器会检查异常名称,以及是否有计数器曾经是 15 的快照。

重新读取,重新计算

用 /root/ice/conc/retry.py 重做 b 的 +5。每次尝试都重新读取表,在当时读到的值上加 5,尝试带条件的覆盖写入,失败就从头再来。结束后,hits 必须是 20。

库的重试只是重新应用提交,不会替你把以前读到的值再读一遍。如果把用以前读到的 10 算出的 15 原样重写,提交会成功,而 a 的 +5 就消失了——这就是丢失更新。评分器会检查计数器是否恰好是 20。

把并发写入规则变成团队规则

在 /root/ice/conc/report.md 中写出 ## 자동 재시도、## 실패한 커밋의 흔적、## 잃어버린 갱신 三个小节。第二节以数字写入第 4 步找到的孤儿文件数,第三节以数字写入最终的计数器值。

如果向同一张表写入的作业有两个以上,请写出哪些作业可以交给重试,哪些作业必须从读取开始重做,以及失败作业留下的文件由谁、在什么时候来清理。