文件怎么摆放决定查询成本 — 分区与合并
目标
创建用于操作文件系统上分区数据湖的工具 pq.py。一边更换键,一边测量文件数和大小分布,测量作用在分区键上的条件能把打开的文件数减少多少,执行文件合并,确认合并过程中读取方会看到什么,并测量更换键的代价。
为什么重要
即使插入同样的行,文件怎么摆放也会让查询打开的文件数相差几十倍。分区就是按值把目录分开摆放,没被打开的文件成本为 0。不过这个好处只有在按目录名中包含的字段过滤时才会出现。 但如果把常用的字段都放进键里,分区数就会按值的乘积增长,一个文件只剩几行。打开文件的固定开销超过了读取内容的开销,列表的条目数暴涨,压缩也起不了作用。所以选键永远是一种取舍。 小文件堆起来之后,就要执行文件合并。难的不是合并本身,而是合并过程中读取的人会看到什么。通过遍历目录来收集文件的读取方,会把旧文件和新文件都捡走,同一行被统计两次。让读取方去看列表,并且只在最后一次性替换列表,这中间的空档就消失了。 本实验不使用 Parquet。实验镜像里没有 pyarrow,Pod 也不能在运行时安装软件包。行组之类的概念已在理论课时中通过官方文档讲过,这里用分区目录、清单和 JSON Lines 手工搭出同样的结构。 评分器不会相信你写出来的文字。它会在临时工作目录里摆好评分器生成的原始数据,真正运行你的工具,检查清单里记录的大小是否与磁盘上的实际大小一致、行是否被保留、裁剪掉的文件数是否正确。原始数据和目标大小每次运行都会变化。
步骤
- 创建并运行 /root/parts/gen_orders.py,生成 /root/parts/work/raw.jsonl。
- 在 /root/parts/pq.py 中实现
write,生成按一个键切分的数据湖和清单。 - 加上
layout,输出文件数、大小分布和小文件个数。 - 让
write能接收多个键并限制单个文件的行数,生成切得很碎的数据湖。 - 加上
query,用作用在分区键上的条件裁剪要打开的文件。 - 加上
compact,把一个分区内的小文件合并到接近目标大小。 - 加上
--crash=before-swap和--via=glob,确认合并过程中读取方会看到什么。 - 加上
repartition,更换键重新写入,并把代价记在 /root/parts/work/partition_report.md 中。
参考
- 所有工作都在
/root/parts下进行。工作目录是/root/parts/work,原始数据是其中的raw.jsonl。 - 一个数据湖就是工作目录下的一个目录。其中放着分区目录和清单
_manifest.json。 - 分区目录名是
칸=값(占位符依次为字段名、值),如果有多个键,就按键的顺序层层嵌套。例如:day=2026-01-03/region=seoul/。 - 分区文件名以
part-开头、以.jsonl结尾。一行就是原始数据的一行。 - 清单:
{"lake": 이름, "key": [칸 이름], "generation": 정수, "files": [{"path": 레이크 기준 상대 경로, "partition": {칸: 값}, "rows": 정수, "bytes": 정수}]}(占位符依次为数据湖名称、字段名、整数、相对于数据湖的相对路径、字段、值、整数、整数)。files 按 path 升序排列,bytes 必须是磁盘上的实际文件大小。 - 运行契约:
python3 /root/parts/pq.py <명령> <작업폴더> [...](占位符依次为命令、工作目录)。结果以一个 JSON 对象输出到标准输出。成功时退出码为 0,缺少工作目录或清单时为 3,用法错误时为 2,进入故意终止的路径时为 9。 write <작업폴더> --lake=<이름> --key=<칸[,칸]> [--rows=<줄 수>]的响应是{"lake": 이름, "key": [칸], "partitions": 정수, "files": 정수, "rows": 정수, "bytes": 정수}(占位符依次为工作目录、名称、字段列表、行数;名称、字段、整数、整数、整数、整数)。--rows是单个文件最多容纳的行数,不传或为 0 时,每个分区一个文件。如果已经有同名的数据湖,就删除后重新写入。layout <작업폴더> --lake=<이름> --small=<바이트>的响应是{"lake": 이름, "files": 정수, "partitions": 정수, "rows": 정수, "bytes": 정수, "avg_bytes": 정수, "p50_bytes": 정수, "min_bytes": 정수, "max_bytes": 정수, "small_files": 정수}(占位符依次为工作目录、名称、字节数;名称,以及九个整数)。avg_bytes 是总字节数除以文件数的商(舍去小数),p50_bytes 是文件大小按最近秩取的中位数,small_files 是大小小于--small的文件数。query <작업폴더> --lake=<이름> --where=<칸=값[,칸=값]> [--via=manifest|glob]的响应是{"lake": 이름, "via": 문자열, "files_total": 정수, "files_scanned": 정수, "rows": 정수, "amount": 정수}(占位符依次为工作目录、名称、条件列表;名称、字符串、整数、整数、整数、整数)。值按字符串比较。只能用分区键中的字段来裁剪文件——如果条件中有不在键里的字段,这个条件什么也裁剪不了,必须打开文件才知道。--via=glob是忽略清单、直接遍历目录中part-*.jsonl的读取方式。compact <작업폴더> --lake=<이름> --target=<바이트> [--crash=before-swap]的响应是{"lake": 이름, "before_files": 정수, "after_files": 정수, "partitions": 정수, "rows": 정수, "bytes": 정수}(占位符依次为工作目录、名称、字节数;名称、整数、整数、整数、整数、整数)。把一个分区内的文件按 path 顺序首尾相接,当加上下一个文件会超过目标时就在那里截断(一个文件无论多大,都单独成为一个文件)。写完全部新文件之后,再换上清单,然后才删除旧文件。新文件名不能与旧名字重复。--crash=before-swap会在写完新文件、即将换上清单之前,以退出码 9 终止。旧文件和旧清单都保持原样。repartition <작업폴더> --from=<이름> --to=<이름> --key=<칸[,칸]>的响应是{"from": 이름, "to": 이름, "rows_read": 정수, "rows_written": 정수, "files_before": 정수, "files_after": 정수, "partitions_after": 정수, "bytes_read": 정수, "bytes_written": 정수}(占位符依次为工作目录、名称、名称、字段列表;名称、名称、七个整数)。它沿着--from数据湖的清单读取,而不是读原始文件。- 报告 MD 的各节标题是
## 무엇을 어떻게 쪼갰나(韩文,意为“把什么、怎么切分了”)、## 작은 파일 문제(韩文,意为“小文件问题”)、## 묶기(韩文,意为“文件合并”)、## 파티션을 바꾸는 비용(韩文,意为“更换分区的代价”)。 - 官方文档:Parquet Concepts · Parquet Configurations · Parquet File Format · os.replace
- 常见错误:把清单里的 bytes 写成计算值而不是实际大小、把按不在键里的字段裁掉的文件也算作已裁剪、先换上清单再写文件、给合并后的新文件沿用旧名字。
- 想用眼睛看文件是怎么摆放的,可以使用
find /root/parts/work/<레이크> -name 'part-*' | head(占位符为数据湖名称)和du -a。
准备一份有多个分区键候选的原始数据
创建并运行 /root/parts/gen_orders.py,生成 /root/parts/work/raw.jsonl。要求至少 200 行,每一行包含 order_id、day、region、channel 和整数 amount,其中 day 至少有 4 种取值,region 至少有 3 种,channel 至少有 3 种。
值的个数很重要。因为本实验要观察的,就是把三个字段都设为键时,分区会增长到多少个。把取值个数相乘,再除以原始数据的行数,就能提前看出一个文件里会剩几行。要固定随机数种子,这样在更换键来比较的过程中,原始数据才不会变动。
按一个键切分,并留下列表
在 /root/parts/pq.py 中实现 write <작업폴더> --lake=<이름> --key=<칸>(占位符依次为工作目录、名称、字段),为每个值创建 칸=값 目录(占位符依次为字段、值),在其下写入 part-0000.jsonl,并在数据湖里留下 _manifest.json。
清单里的 bytes 必须是磁盘上的实际文件大小。如果写的是把各行长度相加算出的值,就会因换行或编码而出现偏差,这个偏差之后会在文件合并时让截断发生在错误的位置。写完文件之后要重新测量大小。清单也先用临时名字写入再改名替换,这样就不会读到写到一半的列表。
测量文件是以什么大小摆放的
加上 layout <작업폴더> --lake=<이름> --small=<바이트>(占位符依次为工作目录、名称、字节数),输出文件数、分区数、行数、总字节数,平均、中位数、最小、最大大小,以及小于 --small 的文件数。
平均值会说谎。一个大文件和几百个小文件混在一起,平均值看上去完好无损。中位数按最近秩取,不要插值。avg_bytes 是总字节数除以文件数的商(舍去小数)。这些值全都是从清单中读取并计算出来的。
切得很碎,看看小文件是怎么来的
让 write 能像 --key=day,region,channel 这样接收多个键,并能用 --rows=<줄 수>(占位符为行数)限制单个文件的最大行数。然后生成切得很碎的数据湖,用 layout 与前面的数据湖做比较。
每多加一个键,分区数就要乘以该字段的取值个数。目录按键的顺序层层嵌套——就是 day=2026-01-03/region=seoul/channel=app/。--rows 不传或为 0 时,每个分区一个文件。把两个数据湖的 small_files 并排放在一起,就能用数字看出失去了什么。
没被打开的文件,成本为 0
加上 query <작업폴더> --lake=<이름> --where=<칸=값[,칸=값]>(占位符依次为工作目录、名称、条件列表),输出符合条件的行数和金额,但要只用分区键中的字段裁剪要打开的文件。响应中同时包含 files_total 和实际打开的 files_scanned。
不在键里的字段没有出现在目录名中,所以用这个条件什么也裁剪不了。这时 files_scanned 必须等于 files_total。如实地统计这一点,就是这一步的全部——在这里虚报,之后就永远找不到到底什么慢了。值按字符串比较。
合并小文件
加上 compact <작업폴더> --lake=<이름> --target=<바이트>(占位符依次为工作目录、名称、字节数),把一个分区内的文件按 path 顺序首尾相接,当加上下一个文件会超过目标时就截断。写完全部新文件之后再换上清单,然后才删除旧文件。
顺序很重要。如果先换上清单,读取方就会去打开还不存在的文件。新文件名不能与旧名字重复——一旦重复,就会覆盖正在读取的文件,该分区整个变空。把代次编号放进名字里,就不会发生重名。合并前后的行数必须相同。
合并过程中读取的人看到什么
给 compact 加上 --crash=before-swap,让它在写完新文件、即将换上清单之前以退出码 9 终止;再给 query 加上 --via=glob,做出忽略清单、直接遍历目录中 part-*.jsonl 的读取方式。终止之后,比较两种读取的答案。
这是本实验的核心场景。通过清单读取的一方看不到任何变化,遍历目录的一方则会把同一行统计两次。终止之后,旧文件和旧清单必须保持原样。之后正常完成文件合并,旧文件就会被删除,两种读取又变得一致。
给更换键标个价
加上 repartition <작업폴더> --from=<이름> --to=<이름> --key=<칸[,칸]>(占位符依次为工作目录、名称、名称、字段列表),不是读原始数据,而是沿着 --from 数据湖的清单全部读出,再按新键重新写入。然后在 /root/parts/work/partition_report.md 中写成四节,标题是 ## 무엇을 어떻게 쪼갰나(韩文,意为“把什么、怎么切分了”)、## 작은 파일 문제(韩文,意为“小文件问题”)、## 묶기(韩文,意为“文件合并”)、## 파티션을 바꾸는 비용(韩文,意为“更换分区的代价”)。
分区键本身就是目录结构,所以更换它就是把所有行重写一遍。第一步要确认 rows_read 和 rows_written 是否相同,而 bytes_read 和 bytes_written 才是这项操作的实际数值。报告里要用数字写出这些字节数——下次开会有人提出要换键时,需要的就是这个唯一的数字。