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

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

清理 30 次小提交 — 按标签、压缩、过期、孤儿文件的顺序

在 TT Lab 中继续学习

目标

把 30 个小批次分 30 次提交写入,复现流式加载所产生的小文件和快照堆积,然后在需要回去的时间点打上 tag,用 rewrite_data_files 做文件合并,用 expire_snapshots 删除旧快照及其文件,再用 remove_orphan_files 清除没有任何快照指向的文件。记录每一步中文件数和快照数如何变化。

为什么重要

Iceberg 不会自己删除任何东西。每次提交都会堆积新文件和新 metadata,即使做了文件合并,旧文件仍被旧快照指向,依然保留。所以不做清理的表,读取会变慢(因为小文件和清单太多),存储也会不断增大。 清理是三件不同的事。文件合并把小文件重写成大文件,使当前快照更快。过期会删除过旧的快照,以及只被它们指向的文件——能通过时间旅行回到的过去也就相应缩短。孤儿文件清理会删除从未被任何快照指向过的文件(失败作业的痕迹)——要留出足够的时间余量,以免连正在使用的文件也被删掉。 三者都无法撤销。所以顺序很重要。如果有需要回去的时间点,必须在过期之前打上 tag,而孤儿文件清理的时间余量不要缩短。

步骤

  1. 用 /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 每个批次提交一次地写入。
  2. 用 /root/ice/mnt/before.py 把当前的快照数、数据文件数和平均文件大小写入 /root/ice/mnt/out/before.json。
  3. 用 /root/ice/mnt/tag.py(应用 ice-mnt-tag)给第十次提交的快照打上 tag batch10,保留期限为 RETAIN 30 DAYS。
  4. 用 /root/ice/mnt/compact.py(应用 ice-mnt-compact)调用 rewrite_data_files,并把结果写入 /root/ice/mnt/out/compact.json。
  5. 用 /root/ice/mnt/expire.py(应用 ice-mnt-expire),以 retain_last => 1 让比现在更旧的快照过期。
  6. 在 /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,然后真正运行一次。
  7. 用 /root/ice/mnt/after.py 把清理之后的快照数、数据文件数和磁盘上的 Parquet 文件数写入 /root/ice/mnt/out/after.json。
  8. 在 /root/ice/mnt/report.md 中写出 ## 압축、## 만료와 태그、## 고아 파일 三个小节。

参考

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、孤儿文件的时间余量)来安排?如果是用流式每分钟提交一次的表,又有什么不同,也请写出来。