课程首页· 术语表· 速查表· 第 5 课 / 共 6 课

Lesson 05 · 工程

写好一个数据集:分区、大小、演进

一个文件变成三千个文件,体积翻倍、查询慢 50 倍。分区粒度是最容易犯的错。

前四课都在讲一个文件。真实的归档是一个目录——按天/按月切开的成百上千个文件。这一课讲三个必须做对的决定:分区切多细、单文件和 row group 多大、schema 变了怎么办。

分区:把过滤条件变成目录名

所谓 Hive 风格分区,就是把某一列的值编进目录名,让引擎连文件都不用打开就能排除:

warehouse/iot/
├── day=2026-01-01/
│   └── data_0.parquet      ← 这一天的全部数据
├── day=2026-01-02/
│   └── data_0.parquet
└── day=2026-01-03/
    └── data_0.parquet

查询 WHERE day='2026-01-02'
  → 引擎看目录名就知道只需打开一个文件,其余连 footer 都不读

这是比 row group 统计更粗、也更便宜的一层过滤:第 4 课的四层跳过发生在文件内部,分区裁剪发生在文件之外。两者叠加,不冲突。

COPY (SELECT *, strftime(ts, '%Y-%m-%d') AS day
      FROM 'data/iot_sorted.parquet'
      ORDER BY device_id, ts)                -- 分区内仍要排序!
  TO 'warehouse/iot'
  (FORMAT parquet, PARTITION_BY (day), COMPRESSION zstd, ROW_GROUP_SIZE 200000);
-- ** 匹配任意层目录;hive_partitioning 让 day 变成一个可查询的列
SELECT day, count(*), round(avg(value), 3)
FROM read_parquet('warehouse/iot/**/*.parquet', hive_partitioning = true)
WHERE day BETWEEN '2026-01-05' AND '2026-01-08'
GROUP BY 1 ORDER BY 1;

实测:分区切太细会发生什么

同一份 4,320,000 行数据,三种切法:

本机实测 · DuckDB 1.5.5 写入与查询 · 取 3 次最好成绩
布局文件数总体积平均文件整天聚合单设备查询全量聚合
不分区,单文件116.36 MB16.4 MB9.3 ms3.1 ms9.3 ms
按 day 分区1614.16 MB906 KB2.5 ms5.0 ms8.6 ms
按 day + device_id321831.42 MB ⚠10.0 KB ⚠131.2 ms ⚠121.6 ms ⚠192.6 ms ⚠

第三行就是小文件灾难的教科书样本。文件数涨 200 倍,代价是:

分区键的作用是把数据切成"查询通常一次要一整块"的粒度,不是把数据切碎。第二个分区键几乎总是错的。

分区粒度怎么定

一条经验规则,够用:让每个分区里的文件落在 100 MB – 1 GB 量级。反推分区粒度:

数据量级建议分区键理由
每天几 GB 以上day(甚至 hour)单日就够大,按天切正好
每天几十 MB – 几百 MBmonth按天切会产生大量小文件
每天几 MB不分区,或按 year靠文件内的排序 + row group 统计就够了
多租户 SaaStenant_id + 时间租户隔离是业务需求(删除、导出、计费),不只是性能考量

本课样本每天只有 ~900 KB,按天分区其实已经偏细了——它在表里看起来赢,是因为数据集太小、全在本地磁盘上。真实的 Pivot 场景要按实际日增量反推。

单文件与 row group 该多大

官方 Configurations 给的建议是 row group 512 MB – 1 GB、data page 8 KB。但要看清它的前提:

"Since an entire row group might need to be read, we want it to completely fit on one HDFS block."
——因为可能要读整个 row group,我们希望它正好装进一个 HDFS block。

这是 HDFS 时代的取值。今天绝大多数人把 Parquet 放在 S3 / OSS 上,没有 HDFS block 这回事,而单机引擎(DuckDB / Polars)的并行单位就是 row group——row group 太大反而没法并行,也没法跳过(第 4 课那个 5 个 rg 只能跳掉一点点的例子)。

实践取值(本课程建议,非官方规定)
参数取值为什么
单个文件128 MB – 1 GB小于 128 MB 开始有小文件开销;大于 1 GB 不便于并行与重写
row_group_size100k – 1M 行(约 64–256 MB)要同时够小(跳过粒度细、能并行)和够大(字典有效、元数据占比低)
写入器默认值pyarrow 1,048,576 行;DuckDB 122,880 行两者差 8.5 倍,别假设"默认就行"
write_page_index归档文件一律开代价极小,第 4 课的第三层跳过靠它

小文件治理:compaction

即使分区定对了,流式写入也会产生小文件——每小时一批就是每天 24 个文件。标准做法是定期合并(compaction):把一个分区里的碎片读进来、排序、重新写成一两个大文件。

-- 把 day=2026-01-08 下的所有碎片合并成一个排好序的文件
COPY (SELECT * EXCLUDE (day)
      FROM read_parquet('warehouse/iot/day=2026-01-08/*.parquet')
      ORDER BY device_id, ts)
  TO 'warehouse/iot/day=2026-01-08/compacted.parquet'
  (FORMAT parquet, COMPRESSION zstd, ROW_GROUP_SIZE 200000);
-- 确认无误后再删除原碎片(顺序不能反!)

Schema 演进:加了一列怎么办

设备固件升级,上报里多了一个 firmware 字段。老文件没有这列,新文件有。混着读会怎样?

-- 默认:按位置对齐列,schema 不一致时可能报错或错位
SELECT * FROM read_parquet('warehouse/iot/**/*.parquet');

-- 正确姿势:按列名合并,缺的列自动补 NULL
SELECT * FROM read_parquet('warehouse/iot/**/*.parquet', union_by_name = true);

实测:老文件 1000 行(7 列)+ 新文件 1000 行(8 列),union_by_name=true 读出 2000 行,其中 firmware 非空 1000 行——老数据那部分自动补 NULL。

裸 Parquet 目录能安全承受的演进只有这几种:

变更安全吗说明
加一列(可空)✅ 安全配合 union_by_name,老文件补 NULL
删一列⚠️ 读得出新文件没这列,查询要容忍 NULL;下游别写死列序
改列名❌ 危险按名合并会变成"两列各半边数据",等同于加一列 + 删一列
改类型(int→bigint)⚠️ 看引擎宽化通常可以,窄化和 int↔string 会报错。改类型前先在小样本上验证

真正的 schema 演进管理(列重命名、类型变更、字段 ID 追踪)是表格式的职责,不是 Parquet 的。这也是判断"该不该上 Iceberg"的第二个信号。

动手做 · 10 分钟 · 亲手制造并量化小文件灾难
import duckdb, os, glob, shutil, time
con = duckdb.connect()
con.sql("CREATE VIEW src AS SELECT *, strftime(ts,'%Y-%m-%d') AS day "
        "FROM 'data/iot_sorted.parquet'")

for d in ["data/p_day", "data/p_day_dev"]: shutil.rmtree(d, ignore_errors=True)
con.sql("COPY (SELECT * FROM src ORDER BY device_id, ts) TO 'data/p_day' "
        "(FORMAT parquet, PARTITION_BY (day), COMPRESSION zstd, ROW_GROUP_SIZE 200000)")
con.sql("COPY (SELECT * FROM src ORDER BY device_id, ts) TO 'data/p_day_dev' "
        "(FORMAT parquet, PARTITION_BY (day, device_id), COMPRESSION zstd, ROW_GROUP_SIZE 200000)")

for d, g in [("p_day", "data/p_day/*/*.parquet"),
             ("p_day_dev", "data/p_day_dev/*/*/*.parquet")]:
    fs = glob.glob(g); tot = sum(os.path.getsize(f) for f in fs)
    s = time.perf_counter()
    con.sql(f"SELECT avg(value) FROM read_parquet('{g}', hive_partitioning=true) "
            f"WHERE device_id='AHU-05-013'").fetchall()
    print(f"{d:10s} 文件 {len(fs):5d}  总计 {tot/1024/1024:6.2f} MB  "
          f"单设备查询 {(time.perf_counter()-s)*1000:7.1f} ms")

预期结果:p_day 约 16 个文件 / 14 MB / 5 ms,p_day_dev 约 3218 个文件 / 31 MB / 120 ms 上下。数量级对上即可,绝对值随机器不同。

这一课的小收获

分区是文件之外的一层过滤,跟第 4 课的四层跳过叠加。但分区切细的代价是暴涨的:文件数 ×200 → 体积 ×2.2、查询慢 13–50 倍。定粒度的规则是让每个分区的文件落在 100 MB – 1 GB,第二个分区键几乎总是错的。row group 用 100k–1M 行(pyarrow 和 DuckDB 的默认值差 8.5 倍,别信默认)。流式写入必然产生碎片,要定期 compaction,但裸目录的 compaction 不是原子的。schema 演进只有"加可空列 + union_by_name"是真正安全的。后两条越难受,就越接近该上表格式的时候。

自测一下

同事说"按 day 和 device_id 两层分区,这样查单台设备最快"。实测结果却慢了 40 倍。最主要的原因是什么?

随时问我:把你要归档的表的日增量大小和最常见的三个查询告诉我,我直接帮你算分区粒度和 row group 大小,并写出对应的 COPY 语句。这是本课程最实用的一次对话。

延伸阅读(本课推荐先读这一篇) DuckDB — Partitioned Writes——PARTITION_BY 的完整语义与官方对小文件的警告,本课分区语法的依据 · Parquet 官方 Configurations(注意它的 HDFS 前提) · pyarrow Tabular Datasets(write_dataset 的分区与文件大小控制)