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

数据流水线

文件怎么摆放决定查询成本 — 分区与合并

在 TT Lab 中继续学习

目标

创建用于操作文件系统上分区数据湖的工具 pq.py。一边更换键,一边测量文件数和大小分布,测量作用在分区键上的条件能把打开的文件数减少多少,执行文件合并,确认合并过程中读取方会看到什么,并测量更换键的代价。

为什么重要

即使插入同样的行,文件怎么摆放也会让查询打开的文件数相差几十倍。分区就是按值把目录分开摆放,没被打开的文件成本为 0。不过这个好处只有在按目录名中包含的字段过滤时才会出现。 但如果把常用的字段都放进键里,分区数就会按值的乘积增长,一个文件只剩几行。打开文件的固定开销超过了读取内容的开销,列表的条目数暴涨,压缩也起不了作用。所以选键永远是一种取舍。 小文件堆起来之后,就要执行文件合并。难的不是合并本身,而是合并过程中读取的人会看到什么。通过遍历目录来收集文件的读取方,会把旧文件和新文件都捡走,同一行被统计两次。让读取方去看列表,并且只在最后一次性替换列表,这中间的空档就消失了。 本实验不使用 Parquet。实验镜像里没有 pyarrow,Pod 也不能在运行时安装软件包。行组之类的概念已在理论课时中通过官方文档讲过,这里用分区目录、清单和 JSON Lines 手工搭出同样的结构。 评分器不会相信你写出来的文字。它会在临时工作目录里摆好评分器生成的原始数据,真正运行你的工具,检查清单里记录的大小是否与磁盘上的实际大小一致、行是否被保留、裁剪掉的文件数是否正确。原始数据和目标大小每次运行都会变化。

步骤

  1. 创建并运行 /root/parts/gen_orders.py,生成 /root/parts/work/raw.jsonl。
  2. 在 /root/parts/pq.py 中实现 write,生成按一个键切分的数据湖和清单。
  3. 加上 layout,输出文件数、大小分布和小文件个数。
  4. 让 write 能接收多个键并限制单个文件的行数,生成切得很碎的数据湖。
  5. 加上 query,用作用在分区键上的条件裁剪要打开的文件。
  6. 加上 compact,把一个分区内的小文件合并到接近目标大小。
  7. 加上 --crash=before-swap 和 --via=glob,确认合并过程中读取方会看到什么。
  8. 加上 repartition,更换键重新写入,并把代价记在 /root/parts/work/partition_report.md 中。

参考

准备一份有多个分区键候选的原始数据

创建并运行 /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 才是这项操作的实际数值。报告里要用数字写出这些字节数——下次开会有人提出要换键时,需要的就是这个唯一的数字。