清理 30 次小提交 — 按标签、压缩、过期、孤儿文件的顺序
目标
把 30 个小批次分 30 次提交写入,复现流式加载所产生的小文件和快照堆积,然后在需要回去的时间点打上 tag,用 rewrite_data_files 做文件合并,用 expire_snapshots 删除旧快照及其文件,再用 remove_orphan_files 清除没有任何快照指向的文件。记录每一步中文件数和快照数如何变化。
为什么重要
Iceberg 不会自己删除任何东西。每次提交都会堆积新文件和新 metadata,即使做了文件合并,旧文件仍被旧快照指向,依然保留。所以不做清理的表,读取会变慢(因为小文件和清单太多),存储也会不断增大。 清理是三件不同的事。文件合并把小文件重写成大文件,使当前快照更快。过期会删除过旧的快照,以及只被它们指向的文件——能通过时间旅行回到的过去也就相应缩短。孤儿文件清理会删除从未被任何快照指向过的文件(失败作业的痕迹)——要留出足够的时间余量,以免连正在使用的文件也被删掉。 三者都无法撤销。所以顺序很重要。如果有需要回去的时间点,必须在过期之前打上 tag,而孤儿文件清理的时间余量不要缩短。
步骤
- 用 /root/ice/mnt/trickle.py(应用
ice-mnt-trickle)以PARTITIONED BY (days(order_ts))和 format-version 2 创建lake.mnt.orders,并把/data/ice/batches/batch-001.csv到batch-030.csv每个批次提交一次地写入。 - 用 /root/ice/mnt/before.py 把当前的快照数、数据文件数和平均文件大小写入 /root/ice/mnt/out/before.json。
- 用 /root/ice/mnt/tag.py(应用
ice-mnt-tag)给第十次提交的快照打上 tagbatch10,保留期限为RETAIN 30 DAYS。 - 用 /root/ice/mnt/compact.py(应用
ice-mnt-compact)调用rewrite_data_files,并把结果写入 /root/ice/mnt/out/compact.json。 - 用 /root/ice/mnt/expire.py(应用
ice-mnt-expire),以retain_last => 1让比现在更旧的快照过期。 - 在
/root/ice/warehouse/mnt/orders/data中放一个旧孤儿文件(stray-old.parquet,修改时间为四天前)和一个新文件(stray-new.parquet),用 /root/ice/mnt/orphans.py(应用ice-mnt-orphans)先以dry_run运行remove_orphan_files,把结果写入 /root/ice/mnt/out/orphans_dry.txt,然后真正运行一次。 - 用 /root/ice/mnt/after.py 把清理之后的快照数、数据文件数和磁盘上的 Parquet 文件数写入 /root/ice/mnt/out/after.json。
- 在 /root/ice/mnt/report.md 中写出
## 압축、## 만료와 태그、## 고아 파일三个小节。
参考
- 过程这样调用:
spark.sql("CALL lake.system.<이름>(table => 'lake.mnt.orders', …)")(占位符为过程名称)。参数传的是值(字符串、TIMESTAMP 字面量),不是表达式。 remove_orphan_files如果把older_than设得比现在早不到 24 小时,就会被拒绝。这是为了防止把正在写入的文件(尚未提交,所以看上去像孤儿的文件)删掉而出事故。不指定时,默认值是三天前。- 常见错误:在打上 tag 之前就运行过期,导致第十个快照消失——无法挽回。如果弄成了这样,请在执行
DROP TABLE lake.mnt.orders PURGE之后从第 1 步重新开始。 - 官方文档:Maintenance · Spark Procedures — rewrite_data_files · expire_snapshots · remove_orphan_files · Branching and Tagging — retention
30 个小批次,30 次提交
创建 /root/ice/mnt/trickle.py,应用名称为 ice-mnt-trickle,创建 lake.mnt.orders(六个列,PARTITIONED BY (days(order_ts)),'format-version' = '2'),并把 /data/ice/batches/batch-001.csv 到 batch-030.csv 依次逐个 append()。
一个批次就是一次提交、一个快照,每个日期分区会产生一个小文件。评分器会检查第 30 次提交之后的 metadata 中,30 个快照是否全都是加上该批次行数的 append。
清理之前的数字
用 /root/ice/mnt/before.py(pyiceberg),把快照数、数据文件数(total-data-files)和平均文件大小(total-files-size ÷ 文件数,整数除法),按 {"snapshots", "data_files", "avg_file_bytes"} 的格式写入 /root/ice/mnt/out/before.json。
当前快照的摘要里有整张表的累计值(total-*),所以不必读取全部清单,也能知道文件数和大小。如果一个 Parquet 文件只有几 KB,读取时打开文件的成本就会超过读取数据的成本。
删除之前先起名字
创建 /root/ice/mnt/tag.py,应用名称为 ice-mnt-tag,按 committed_at 的顺序读取 lake.mnt.orders.snapshots,给第十个快照打上 CREATE TAG batch10 AS OF VERSION <ID> RETAIN 30 DAYS。
过期会删除比 older_than 更旧的快照,但 tag 或分支所指向的快照,在该 tag 的保留期限还没结束之前不会被删除。所以 tag 必须在过期之前打上。评分器会检查 tag 是否指向第 30 次提交之后 metadata 中的第十个快照。
文件合并——重写小文件
创建 /root/ice/mnt/compact.py,应用名称为 ice-mnt-compact,调用 CALL lake.system.rewrite_data_files(table => 'lake.mnt.orders', options => map('min-input-files', '2')),并把结果行的 rewritten_data_files_count 和 added_data_files_count,按 {"rewritten", "added"} 的格式写入 /root/ice/mnt/out/compact.json。
文件合并会读取同一分区的小文件,把它们重写成大文件,并用一个“replace”快照提交。一行数据也不会改变。旧的小文件只是从列表中去掉,在磁盘上依然存在——因为旧快照(和 tag)还指向它们。评分器会把 replace 提交的摘要与你的两个值对照。
过期——旧快照及其文件被删除
创建 /root/ice/mnt/expire.py,应用名称为 ice-mnt-expire,读取当前时间,然后调用 CALL lake.system.expire_snapshots(table => 'lake.mnt.orders', older_than => TIMESTAMP '<지금>', retain_last => 1)(占位符为当前时间)。
过期会把快照从 metadata 中删除,并把剩下的快照都没有指向的文件从磁盘上删除。现在只剩 main 和 tag batch10 的快照,而 tag 所指向的第十个快照的小文件不会被删除。评分器会检查剩下的快照和磁盘上的文件。
孤儿文件清理——留出时间余量
在 /root/ice/warehouse/mnt/orders/data 中复制一个现有的数据文件,创建 stray-old.parquet(修改时间用 touch -d '4 days ago')和 stray-new.parquet(当前时间)。然后创建 /root/ice/mnt/orphans.py,应用名称为 ice-mnt-orphans,把 remove_orphan_files(table => 'lake.mnt.orders', dry_run => true) 结果中的 orphan_file_location 每行一个写入 /root/ice/mnt/out/orphans_dry.txt,接着不带 dry_run 再调用一次。
两个文件都是没有任何快照指向的孤儿。但是,刚刚产生的文件,可能是现在有人正在写入的某次提交的文件。所以不会动比 older_than 更新的文件(默认三天)。评分器会检查 dry_run 列表中是否只有旧的,以及实际上是否只删除了旧的。
清理之后的数字
用 /root/ice/mnt/after.py(pyiceberg),把剩下的快照数、当前快照的数据文件数,以及 /root/ice/warehouse/mnt/orders/data 之下的 Parquet 文件数,按 {"snapshots", "data_files", "files_on_disk"} 的格式写入 /root/ice/mnt/out/after.json。
磁盘上的文件数比当前快照的文件数多。因为有 tag batch10 所保护的小文件,以及孤儿文件清理有意留下的新文件。如果能解释这个差别,就说明理解了清理作业。
把清理作业变成运维流程
在 /root/ice/mnt/report.md 中写出 ## 압축、## 만료와 태그、## 고아 파일 三个小节。第一节以数字写入文件合并之前的数据文件数(第 2 步)和文件合并之后当前快照的数据文件数(第 7 步)。
如果把这三件事做成每天运行的作业,你会按什么顺序、周期和标准(older_than、retain_last、孤儿文件的时间余量)来安排?如果是用流式每分钟提交一次的表,又有什么不同,也请写出来。