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

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

同样的删除与 MERGE 两种做法 — 重写文件,还是记下要删的行

在 TT Lab 中继续学习

目标

用同样的三月订单创建 copy-on-write(COW)表和 merge-on-read(MOR)表,对两张表运行同样的 DELETE 和同样的 MERGE。通过快照摘要和文件确认:COW 会把含有被修改行的数据文件整个重写,而 MOR 只增加记录“删除某个文件的第几行”的位置删除文件。亲手打开位置删除文件,并看其他引擎(DuckDB、pyiceberg)读取时是否反映了这些删除。

为什么重要

Parquet 文件无法修改。要改一行,要么把那个文件重新写一遍(COW),要么把“这一行已被删除”这件事记在另一个文件里,读取时再过滤掉(MOR)。COW 读取便宜、写入昂贵——在一百万行的文件里只改一行,也要重写一百万行。MOR 写入便宜、读取昂贵——每次读取都要把删除文件与数据文件对一遍,而且删除文件越积越多,就越慢。 像 CDC 这样频繁出现小变更的表,通常用 MOR 写入并定期做文件合并;每天集中大改一次、读取很多的表,COW 更合适。无论哪种,都由三个表属性(write.delete.mode、write.update.mode、write.merge.mode)决定,而判断的依据不是感觉,而是每次提交都会留在摘要中的写入字节数和文件数。

步骤

  1. 用 /root/ice/rl/tables.py(应用 ice-rl-tables)以 format-version 2 创建 lake.rl.cow(三个模式都是 copy-on-write)和 lake.rl.mor(三个模式都是 merge-on-read),并分别一次性写入三月整月的数据。
  2. 用 /root/ice/rl/delete_cow.py(应用 ice-rl-delete-cow)运行 DELETE FROM lake.rl.cow WHERE status = 'cancelled'。
  3. 用 /root/ice/rl/delete_mor.py(应用 ice-rl-delete-mor)对 lake.rl.mor 运行同样的删除。
  4. 用 /root/ice/rl/merge.py(应用 ice-rl-merge)把 /data/ice/changes.csv(op 为 U、D、I)对两张表做同样的 MERGE。
  5. 用 /root/ice/rl/posdel.py 以 pyarrow 打开 lake.rl.mor 的一个位置删除文件,写入 /root/ice/rl/out/posdel.json。
  6. 用 /root/ice/rl/cost.py 把两张表的 MERGE 提交写入的字节数(added-files-size)写入 /root/ice/rl/out/cost.json。
  7. 用 /root/ice/rl/read.py 分别用 DuckDB 和 pyiceberg 读取 lake.rl.mor,写入 /root/ice/rl/out/read.json。
  8. 在 /root/ice/rl/report.md 中写出 ## 삭제 두 방식、## 위치 삭제 파일、## 쓰기와 읽기의 맞바꿈 三个小节。

参考

同样的数据,不同的写入模式

创建 /root/ice/rl/tables.py,应用名称为 ice-rl-tables,创建 lake.rl.cow 和 lake.rl.mor。两张表都是六个列、'format-version' = '2',并把 write.delete.mode、write.update.mode、write.merge.mode,cow 全部设为 copy-on-write,mor 全部设为 merge-on-read。把三月的 31 个文件分别写入每张表一次。

写入模式是表属性,所以记住它的是表,而不是引擎。目的是让无论哪个引擎来写,都遵循同样的方式。评分器会检查六个属性和第一次提交的行数。

COW 删除——重写文件

创建 /root/ice/rl/delete_cow.py,应用名称为 ice-rl-delete-cow,并运行 DELETE FROM lake.rl.cow WHERE status = 'cancelled'。

含有已取消订单的数据文件,会换成去掉已取消订单的新文件。摘要中的 deleted-data-files、added-data-files 就是这个,而且一个删除文件都没有。评分器会检查 cow 的第二次提交摘要,以及那时是否没有残留已取消订单。

MOR 删除——把要删除的行记下来

创建 /root/ice/rl/delete_mor.py,应用名称为 ice-rl-delete-mor,并运行 DELETE FROM lake.rl.mor WHERE status = 'cancelled'。

数据文件保持原样,只增加一个汇集了要删除行的(文件路径,行号)的位置删除文件。摘要中会出现 added-position-delete-files,而没有 deleted-data-files。评分器会检查 mor 的第二次提交摘要,以及清单中 content 为 1 的文件。

MERGE——修改、删除、添加

创建 /root/ice/rl/merge.py,应用名称为 ice-rl-merge,把 /data/ice/changes.csv 读取为临时视图,对两张表分别用 MERGE INTO … ON t.order_id = s.order_id,对 op 为 D 的执行 DELETE,对 op 为 U 的 UPDATE status 和 amount,对表中没有的 op 为 I 的执行 INSERT。

如果源中的一行匹配了两个变更,MERGE 就会失败(这份变更集中每个订单只有一个变更)。针对已被删除的已取消订单的 U、D,没有匹配的行,什么都不会做。评分器会把两张表的内容与按原始数据和变更集算出的期望值对照,并检查位置删除文件是否只存在于 mor 中。

打开位置删除文件看一看

用 /root/ice/rl/posdel.py,在 lake.rl.mor 的当前快照中,用 pyarrow.parquet.read_table 打开一个 content 为 1 的文件,并把 {"delete_file": 경로, "rows": 행 수, "targets": [그 파일이 가리키는 데이터 파일 경로, …]} 写入 /root/ice/rl/out/posdel.json(占位符依次为删除文件路径、行数与该文件所指向的数据文件路径)。

位置删除文件是只有 file_path、pos 两列的 Parquet。一行的意思是“这个数据文件的第 pos 行不存在”。评分器会把你写下的行数和目标列表,重新从文件里读取来对照,并检查目标是否是当前存活的数据文件。

一次 MERGE 写入的字节数

用 /root/ice/rl/cost.py,从两张表当前快照(= MERGE 提交)的摘要中读取 added-files-size,按 {"cow_added_bytes": 정수, "mor_added_bytes": 정수} 的格式写入 /root/ice/rl/out/cost.json(占位符均为整数)。

COW 会把变更所涉及的数据文件全部重写,MOR 只写新行、被修改的行和删除文件。同样是 1,400 条变更,写入字节数的差距就是写放大。评分器会把两个值与摘要对照,并检查 COW 一侧是否更大。

其他引擎读取时是否也反映了删除

用 /root/ice/rl/read.py,把 lake.rl.mor 当前的 metadata 路径传给 iceberg_scan(),用 DuckDB 统计行数和 sum(amount),再用 pyiceberg 统计行数,按 {"duckdb_rows", "duckdb_amount", "pyiceberg_rows"} 的格式写入 /root/ice/rl/out/read.json。

MOR 表只有读取的一方应用了删除文件,才能得到正确的结果。不认识删除文件的引擎会把已删除的行也返回来——如果多个引擎读取同一张表,这是必须确认的事情。DuckDB 的 iceberg 扩展已在构建时放进了镜像(LOAD iceberg)。评分器会把这三个值与期望值对照。

哪张表用哪种模式

在 /root/ice/rl/report.md 中写出 ## 삭제 두 방식、## 위치 삭제 파일、## 쓰기와 읽기의 맞바꿈 三个小节。第三节以数字写入第 6 步的两个字节值。

想一想你们团队的两张表,如果一张定为 COW,一张定为 MOR,依据是什么?如果选了 MOR,还请写出,堆积的删除文件将由谁、在什么时候来清理。