分区是规则而不是列 — 隐藏分区与分区演进
一句话总结
Iceberg 的分区值不是人工创建的列,而是由针对源列设置的转换规则(例如 days(order_ts))计算出来并写进清单,读取的一方只用源列条件就能跳过文件。规则带着 partition spec 的 ID 累积下来,所以中途改变规则,旧文件仍按旧规则原样读取。
为什么 Hive 风格的分区是个问题
Partitioning 文档举的例子很典型。要按日期划分一张日志表,在 Hive 里得另外创建一个名为 event_date 的列,由写入一方从 event_time 计算出这个值再填进去。问题就是从这里一连串冒出来的。
- 写入一方出错,就会悄悄地错。把
2018-12-01写成20181201,用错了源列(处理时间),或者时区算错了,都不会报错。只是结果错了。 - 读取一方必须知道这一点才能变快。只对
event_time加条件的话,Hive 不知道这两列之间的关系,会读取所有文件。人们必须记住要另外加上event_date条件。 - 无法更改。查询被绑定在分区列上,想把按天改成按小时,就得新建一张表,并把所有查询都改一遍。
工作原理——转换与 partition spec
Iceberg 的 partition spec(分区规范),每个字段都有源列 ID、分区字段 ID、转换和名称。引擎写入行时会计算转换,并把它作为该文件的分区值写进清单。一个文件中的所有行都具有相同的分区值。转换列表如下。
| 转换 | 结果 | 用途 |
|---|---|---|
| identity | 原值 | 基数低的类别列 |
| year · month · day · hour | 自 1970 年起计数的年、月、日、小时 | 时间列 |
| bucket[N] | 哈希值对 N 取余 | 基数高的键(客户编号) |
| truncate[W] | 按宽度 W 截断的值 | 数值区间、字符串前缀 |
| void | 始终为 null | 在版本 1 中去掉字段时 |
bucket 是丢弃 32 位 Murmur3(x86,种子 0)哈希值的符号位后,再对 N 取余。因为规范规定了哈希函数,所以即使由 Spark 写入、由 Python 读取,相同的值也会进入同一个桶。
读取的一方只使用 order_ts >= X 这样的源列条件。扫描规划会把这个条件转换成分区条件(order_ts_day >= day(X)),再与清单中的分区值对比。这种转换是按“包含”的方向计算的,所以可能含有满足条件的行的文件绝不会被漏掉。用户不必知道分区是怎么划分的——所以叫“隐藏”分区。
分区演进——旧文件原样不动
假设数据增多,按月划分变得太粗了。根据 Evolution 文档,即使更改了 spec,旧数据仍然使用旧 spec,只有新数据才按新的布局写入。spec 会累积在列表中,每个清单都会记住自己所用的 spec ID。规划时,会针对每个 spec 分别转换条件来选取文件,文档称之为 split planning。文档还明确写道,分区演进是 metadata 操作,不会急着重写文件。
在 Spark 中,可以用 ALTER TABLE 添加字段(ADD PARTITION FIELD)、删除字段(DROP PARTITION FIELD)、替换字段(REPLACE PARTITION FIELD … WITH …)。
ALTER TABLE lake.demo.events REPLACE PARTITION FIELD months(event_ts) WITH days(event_ts);
在现场相遇的样子
我改了分区,所以旧数据也得重写。这是最常见的误解。旧的按月文件即使原样不动,也能被准确读取。重写只是一种选择,而且必须有理由——比如经常按天查询旧的时间段,想要减少读取量。即便那时,也要通过文件合并(rewrite_data_files)这样的单独作业来做。
条件只是一天,却要读一个月的数据。对按月的文件设置按天的条件,在分区层面就得整体选中那个月的文件。文件内部的列统计信息(下界、上界)会有所帮助,但如果文件均匀地包含整个月,就没有用。规划所打开文件的行数,与实际命中的行数之间的差距,就是这部分成本。
按客户编号划分,结果有几万个分区。用 identity 去划分基数高的列,会导致小文件爆炸。应该用 bucket 把桶数固定下来。
实际工作中真正重要的事
- 分区是规则。只写源列和转换,分区值由引擎来计算。
- 查询只针对源列设置条件。不需要另外知道分区列。
- spec 会累积,文件记得自己的 spec。演进之后,不必重写旧文件。
- 基数高的列用 bucket,时间用 year · month · day · hour。
下一项实验要做什么
把三月整月的数据放进按 months(order_ts) 划分的表,用 pyiceberg 针对一天的条件做规划,测出要打开多少个文件、多少行。把分区改成 days(order_ts) 之后写入四月,确认三月的文件仍是 spec 0,只有四月的文件是 spec 1。对比三月某一天和四月某一天的条件所读取的行数,再把客户表按 bucket(4, customer_id) 划分,检查桶是否按规范计算。