添加、重命名、放宽、删除再添加列 — 一个文件都不重写
目标
用四种方式修改 Iceberg 表的 schema(添加列、重命名、拓宽类型、删除后重新添加),确认旧数据文件一次也没有被重写,却依然能被正确读取。亲手打开文件,看一看其中的奥秘:它就是写在 Parquet 文件页脚里的 field_id。
为什么重要
在按名称查找列的表(大部分 Hive 风格的 Parquet 表)中,重命名就是事故。旧文件里只有旧名称,用新名称读取,这一列会整个变成 null。在按位置查找列的格式(CSV)中,一调整顺序,值就会跑进错误的列。所以人们每次修改 schema,都要把整张表重写一遍,或者干脆改不了,带着名称有误的列过好几年。 Iceberg 给每一列分配一个永不改变的 ID,写文件时也把这个 ID 一并写入。读取时不是按名称,而是按 ID 配对。所以重命名只是 metadata 的一行,而且即使用被删除列的名称添加新列,它也是新的 ID,旧值不会复活。类型只能朝值不会被截断的方向(int → long 这样的拓宽)修改。
步骤
- 用 /root/ice/sch/base.py(应用
ice-sch-base)创建lake.sch.orders(format-version 2)并写入 2026-03-01。 - 用 /root/ice/sch/add.py(应用
ice-sch-add)添加coupon STRING列,然后写入 2026-03-02。当天各行的coupon,在amount >= 100000时为'SPRING',否则为 null。 - 用 /root/ice/sch/rename.py(应用
ice-sch-rename)把amount重命名为amount_krw。 - 用 /root/ice/sch/footer.py(pyarrow、pyiceberg)打开第一次提交的一个数据文件,把字段 ID 4 在文件中的名称写入 /root/ice/sch/out/footer.json。
- 用 /root/ice/sch/widen.py(应用
ice-sch-widen)把amount_krw拓宽为BIGINT,并把再尝试收窄为INT时的错误条件名称写入 /root/ice/sch/out/narrow.txt。 - 用 /root/ice/sch/readd.py 删除
coupon,再用同样的名称重新添加,把旧 ID、新 ID 和现在非 null 值的数量写入 /root/ice/sch/out/readd.json。 - 用 /root/ice/sch/history.py 把 schema 数、当前 schema ID、快照数和存活的数据文件数写入 /root/ice/sch/out/history.json。
- 在 /root/ice/sch/report.md 中写出
## 이름 바꾸기、## 형 넓히기、## 지웠다 다시 더하기三个小节。
参考
- Spark 的
df.schema中没有字段 ID。ID 要通过 metadata.json 的schemas[].fields[].id或 pyiceberg 的tbl.schema()来看。 - Parquet 页脚用
pyarrow.parquet.read_schema(경로)读取(占位符为文件路径),每个字段的metadata[b"PARQUET:field_id"]就是 Iceberg 字段 ID。 - 修改 schema 不会产生快照。这个实验结束后,快照应当只有两个(3 月 1 日和 2 日)。
- 常见错误:重新运行第 1、2 步的脚本,使快照变成三个。如果弄成了这样,请在执行
spark-sql -e "DROP TABLE lake.sch.orders PURGE"之后从第 1 步重新开始。 - 官方文档:Evolution — Schema evolution · Spec — Schema Evolution · Spark DDL — ALTER TABLE
表与第一次提交——每一列都有 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 读取的使用方(直接打开文件的脚本),还要写出会有什么被破坏。