同样的删除与 MERGE 两种做法 — 重写文件,还是记下要删的行
目标
用同样的三月订单创建 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)决定,而判断的依据不是感觉,而是每次提交都会留在摘要中的写入字节数和文件数。
步骤
- 用 /root/ice/rl/tables.py(应用
ice-rl-tables)以 format-version 2 创建lake.rl.cow(三个模式都是copy-on-write)和lake.rl.mor(三个模式都是merge-on-read),并分别一次性写入三月整月的数据。 - 用 /root/ice/rl/delete_cow.py(应用
ice-rl-delete-cow)运行DELETE FROM lake.rl.cow WHERE status = 'cancelled'。 - 用 /root/ice/rl/delete_mor.py(应用
ice-rl-delete-mor)对lake.rl.mor运行同样的删除。 - 用 /root/ice/rl/merge.py(应用
ice-rl-merge)把/data/ice/changes.csv(op 为 U、D、I)对两张表做同样的 MERGE。 - 用 /root/ice/rl/posdel.py 以 pyarrow 打开
lake.rl.mor的一个位置删除文件,写入 /root/ice/rl/out/posdel.json。 - 用 /root/ice/rl/cost.py 把两张表的 MERGE 提交写入的字节数(
added-files-size)写入 /root/ice/rl/out/cost.json。 - 用 /root/ice/rl/read.py 分别用 DuckDB 和 pyiceberg 读取
lake.rl.mor,写入 /root/ice/rl/out/read.json。 - 在 /root/ice/rl/report.md 中写出
## 삭제 두 방식、## 위치 삭제 파일、## 쓰기와 읽기의 맞바꿈三个小节。
参考
- 变更集的列是
op, order_id, customer_id, region, amount, status, order_ts。U 把状态改为 refunded,D 表示删除,I 表示新订单。 - 提交摘要用
SELECT operation, summary FROM lake.rl.mor.snapshots查看,文件类型用SELECT content, file_path, record_count FROM lake.rl.mor.files(content 0 数据 · 1 位置删除 · 2 等值删除)查看。 - 这个实验的表是 format-version 2,所以位置删除文件是以 Parquet 写成的。在 format-version 3 中,同样的信息会以 deletion vector(Puffin)写成。
- 常见错误:把第 2 步对两张表都运行一遍(与第 3 步的顺序混在一起,对比就会错位),以及把 MERGE 运行两次。想恢复,请在对两张表执行
DROP TABLE … PURGE之后,从第 1 步重新开始。 - 官方文档:Spark Writes — MERGE INTO · Configuration — write.delete.mode · Spec — Row-level Deletes · DuckDB — Iceberg extension
同样的数据,不同的写入模式
创建 /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,还请写出,堆积的删除文件将由谁、在什么时候来清理。