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

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

添加、重命名、放宽、删除再添加列 — 一个文件都不重写

在 TT Lab 中继续学习

目标

用四种方式修改 Iceberg 表的 schema(添加列、重命名、拓宽类型、删除后重新添加),确认旧数据文件一次也没有被重写,却依然能被正确读取。亲手打开文件,看一看其中的奥秘:它就是写在 Parquet 文件页脚里的 field_id。

为什么重要

在按名称查找列的表(大部分 Hive 风格的 Parquet 表)中,重命名就是事故。旧文件里只有旧名称,用新名称读取,这一列会整个变成 null。在按位置查找列的格式(CSV)中,一调整顺序,值就会跑进错误的列。所以人们每次修改 schema,都要把整张表重写一遍,或者干脆改不了,带着名称有误的列过好几年。 Iceberg 给每一列分配一个永不改变的 ID,写文件时也把这个 ID 一并写入。读取时不是按名称,而是按 ID 配对。所以重命名只是 metadata 的一行,而且即使用被删除列的名称添加新列,它也是新的 ID,旧值不会复活。类型只能朝值不会被截断的方向(int → long 这样的拓宽)修改。

步骤

  1. 用 /root/ice/sch/base.py(应用 ice-sch-base)创建 lake.sch.orders(format-version 2)并写入 2026-03-01。
  2. 用 /root/ice/sch/add.py(应用 ice-sch-add)添加 coupon STRING 列,然后写入 2026-03-02。当天各行的 coupon,在 amount >= 100000 时为 'SPRING',否则为 null。
  3. 用 /root/ice/sch/rename.py(应用 ice-sch-rename)把 amount 重命名为 amount_krw。
  4. 用 /root/ice/sch/footer.py(pyarrow、pyiceberg)打开第一次提交的一个数据文件,把字段 ID 4 在文件中的名称写入 /root/ice/sch/out/footer.json。
  5. 用 /root/ice/sch/widen.py(应用 ice-sch-widen)把 amount_krw 拓宽为 BIGINT,并把再尝试收窄为 INT 时的错误条件名称写入 /root/ice/sch/out/narrow.txt。
  6. 用 /root/ice/sch/readd.py 删除 coupon,再用同样的名称重新添加,把旧 ID、新 ID 和现在非 null 值的数量写入 /root/ice/sch/out/readd.json。
  7. 用 /root/ice/sch/history.py 把 schema 数、当前 schema ID、快照数和存活的数据文件数写入 /root/ice/sch/out/history.json。
  8. 在 /root/ice/sch/report.md 中写出 ## 이름 바꾸기、## 형 넓히기、## 지웠다 다시 더하기 三个小节。

参考

表与第一次提交——每一列都有 ID

创建 /root/ice/sch/base.py,应用名称为 ice-sch-base,创建 lake.sch 和 lake.sch.orders(六个列,'format-version' = '2'),并写入 2026-03-01 的文件。

创建表时,Iceberg 会从 1 开始为每一列编号(order_id 为 1 …… order_ts 为 6)。这个 ID 不会因为名称改变而改变。评分器会从 metadata 的 schema 中查看 ID 和名称,从第一个快照中查看行数。

添加列——旧文件读出来是 null

创建 /root/ice/sch/add.py,应用名称为 ice-sch-add,先执行 ALTER TABLE lake.sch.orders ADD COLUMN coupon STRING,再给 2026-03-02 的文件加上 coupon(amount >= 100000 时为 'SPRING',否则为 null)后写入。

新列会得到新的 ID(7)。3 月 1 日的文件里没有 ID 7,所以那些行的 coupon 读出来是 null——因为没有重写文件。评分器会直接打开 3 月 2 日提交所写的 Parquet 文件,检查 ID 7 这一列的值是否符合原始条件。

重命名——metadata 的一行

创建 /root/ice/sch/rename.py,应用名称为 ice-sch-rename,并运行 ALTER TABLE lake.sch.orders RENAME COLUMN amount TO amount_krw。

重命名是往 metadata 中添加一个新 schema(只有 ID 4 的名称不同)。两天的文件原封不动,sum(amount_krw) 却能得出两天的总和。评分器会检查 ID 4 的名称是否已改变,以及用新名称读出的总和是否与原始总和相同。

Parquet 页脚里留下的旧名称

用 /root/ice/sch/footer.py,通过 pyarrow.parquet.read_schema 打开第一次提交(3 月 1 日)的一个数据文件,找出 PARQUET:field_id 为 4 的字段在文件中的名称,并按 {"file", "name_in_file", "field_id", "name_in_table"} 的格式写入 /root/ice/sch/out/footer.json。

文件原样保存着写入那一刻的名称(amount)。表使用的是现在的名称(amount_krw)。把两者连起来的,就是页脚里的 field_id。第一次提交的文件,可以用 pyiceberg 的 tbl.inspect.files(첫_스냅샷_ID) 找到(占位符为第一个快照的 ID;路径前面的 file: 请去掉)。

类型只能拓宽

创建 /root/ice/sch/widen.py,应用名称为 ice-sch-widen,把 amount_krw 改为 BIGINT,然后把试图改回 INT 的语句用 try 包起来,把异常的 getCondition() 写入 /root/ice/sch/out/narrow.txt 的第一行。

int → long 不会截断任何值,所以旧文件(以 int 写入的)可以原样按 long 读取。反过来值可能被截断,规范不允许。评分器会检查 ID 4 的类型是否为 long,以及错误条件的名称。

删除后用同样的名称重新添加会怎样

用 /root/ice/sch/readd.py 对 coupon 执行 DROP COLUMN,再用同样的名称执行 ADD COLUMN coupon STRING,并把旧 ID、新 ID 和现在的 count(coupon),按 {"old_id", "new_id", "non_null"} 的格式写入 /root/ice/sch/out/readd.json。

被删除的列的 ID 不会再被使用。新的 coupon 会得到新的 ID,3 月 2 日文件中残留的 ID 7 的“SPRING”值与新列配不上,所以看不到。如果是按名称读取的表,旧值就会复活。ID 可以通过 pyiceberg 的 tbl.schemas()(历史)和 tbl.schema()(当前)读取。

五个 schema,两个快照,文件原样不动

用 /root/ice/sch/history.py(pyiceberg),把 schema 数、当前 schema ID、快照数和存活的数据文件数,按 {"schemas", "current_schema_id", "snapshots", "data_files"} 的格式写入 /root/ice/sch/out/history.json。

每次修改 schema,metadata 中就会多累积一个 schema,但快照不会增加,数据文件也保持两次提交所写入的样子。评分器会把这四个值与 metadata 对照,同时也检查现在存活的文件是否与第二个快照的文件完全相同(有没有被重写的文件)。

把 schema 变更规则变成团队规则

在 /root/ice/sch/report.md 中写出 ## 이름 바꾸기、## 형 넓히기、## 지웠다 다시 더하기 三个小节。第三节以数字写入第 6 步的旧 ID 和新 ID。

如果有多个团队读取这张表,请写出哪些变更可以随时进行,哪些变更需要通知。如果有不按字段 ID 读取的使用方(直接打开文件的脚本),还要写出会有什么被破坏。