2. 核心细节解析与实操要点
2.1 嵌套数据的拍平逻辑与控制参数
在动手写代码之前,先把我理解的 Dremel 思路讲透,否则你会在字段展开和 null 处理上被折磨到怀疑人生。Dremel 的论文里定义了 record 和 column 两种视角,核心是把一棵嵌套 JSON 树通过repeated字段展开成多行。对应到 Parquet 物理层,就是 record shredding 和 assembly 的过程。你要是不把repeated和required/optional的语义想清楚,写出来的 schema 和真实数据会有“对不上”的错位感。
我把一个典型订单数据拍平前后做了对比,你可以直观感受一下什么叫“列式存储下的嵌套展开”:
| 字段路径 | 原始 JSON 形态 | 拍平后的列名 | 重复次数 | 说明 |
|---|---|---|---|---|
| order_id | 普通字符串 | order_id | 1 | 每条记录对应一个值 |
| items | 数组,元素含 product 和 price | items.product, items.price | 随数组长度变化 | 需用 repeated 修饰 |
| customer.name | 嵌套对象 | customer.name | 1 | 合并路径即可 |
| tags | 字符串数组 | tags | 随数组长度变化 | 数组本身就是 repeated |
有了这张表,你回头再看encoding里的RLE/bit-packing、DELTA_BINARY_PACKED这些名字,就不会觉得它们是抽象概念了——它们都在为你这种“列内需要记录重复次数”的场景服务。Dremel 当年提出 repetition level 和 definition level,本质上就是给“一个值到底嵌套在哪一层”做了精细标注,Parquet 把这两个概念原封不动继承了下来。
2.2 Parquet 压缩策略与存储格式的选择
如果你只是把 CSV 转成 Parquet 就收工,那我建议你最好再看一眼压缩算法的选择和 row group 大小的设置。我实测下来,在数据量 10GB 左右、单条记录字段数 50 上下时,压缩算法不是越“重”越好。快速给一个对比结论:
snappy:解压速度极快,压缩率中等,适合 OLAP 场景,大多数查询引擎默认用它。gzip:压缩率最高但 CPU 消耗大,适合冷数据归档,不适合高频查询。zstd:压缩率和速度都均衡,新版 Spark、Trino 支持很好,如果你在构建新链路,我优先推荐它。
你还要注意row_group_size。默认值在不同引擎里不一样,但有些开源组件默认是 128MB 或 256MB,我通常设成 512MB,这样既能减少元数据开销,又能让谓词下推更好地命中 row group 内的统计信息。这里背后逻辑是:Parquet 在读取时会用 row group 的 min/max 值做过滤,你把 row group 调大了,每组的统计信息覆盖范围就大,但对查询下推来说反而更容易裁剪掉无关数据块。
- 分区策略和排序键:如果你要按订单日期查询,建议先对日期字段做 partition,再用明细字段做 sort order。分区目录虽然会碎,但是裁剪效率能提升一个量级。
- 字段重排:把高基数字段和低基数字段分开排列,能显著提升压缩率。我是把维度字段放在前,度量字段放在后,中间夹上日期类字段。
- Null 处理:Parquet 的列式存储对 null 有专门编码位图,但大量 null 会让数据膨胀。所以数据清洗时尽量把空字符串转成真正的 null,而不是让“空字符串”占一个字典位置。
2.3 给一个 Python 存取 Parquet 数据的案例
你如果是 Python 生态的重度用户,那我必须给你一个直接用pyarrow读写 Parquet 的案例。别绕去用 pandas 的to_parquet,虽然它内部也是调 pyarrow,但绕一层会导致你对底层控制变弱。直接上代码,最基础的一条路径:
import pyarrow as pa import pyarrow.parquet as pq data = { "order_id": [1001, 1002, 1003], "customer_name": ["Alice", "Bob", "Carol"], "amount": [250.0, 178.5, 920.0], "tags": [ ["new", "vip"], ["return"], ["vip", "promo", "friend"] ], } table = pa.table(data) pq.write_table(table, "orders.parquet") table_read = pq.read_table("orders.parquet") print(table_read.to_pandas())这段代码看起来平铺直叙,但你注意到没有:tags是列表类型,pyarrow 会自动推断成一个list<item: string>的嵌套列,它在 Parquet 里会被映射成 repeated 字段。如果你用 pandas 直接存,很多时候 list 会被解释成 object 列,后续查询引擎没法做很好的下推优化。所以我的习惯是:先用 pyarrow 构造好原始 table,再统一落盘。
还有几个参数我在生产中常用,一次性给你:
pq.write_table( table, "orders_zstd.parquet", compression="zstd", row_group_size=512 * 1024 * 1024, # 512MB version="2.6", data_page_version="1.0", )这里版本号别乱填。Parquet 2.6 支持更新的逻辑类型,比如TIMESTAMP_NANOS,如果你只要普通时间戳,可以保持默认;但如果你要支持微秒或纳秒精度,用version="2.6"会稳妥很多。我在实际业务里遇到过因为版本过低导致时间精度被截断的坑,所以这里多嘴提醒一句。
3. 实操过程与核心环节实现
3.1 数据建模:从 JSON 到 Parquet Schema
在动手写 Schema 之前,你先想清楚一个问题:你的数据更偏向“宽表”还是“嵌套结构”?如果是宽表场景,Parquet 的优势在于压缩和裁剪,字段多的情况下列式存储极其友好;如果是嵌套结构,比如订单含明细行、商品含多规格,你就要合理设计 repeated 层级,否则会产生大量冗余的 repetition level 标记,让文件膨胀。
我给一个实际中验证过的建模流程,前提是数据来自 Kafka JSON 消息,每天千万级:
- 先用 pyarrow 读取几万条 JSON 样例,让 Arrow 自动推断出 schema。注意:这个推断结果往往偏保守,比如数值会被推断成 int64 或 double,字符串会被推断成 string。你需要人工干预,把日期类字段显式转成
timestamp[ns],把金额类字段固定成decimal(18, 2),避免后续查询引擎类型不一致导致谓词下推失效。 - 把嵌套结构拍平:如果嵌套的深度超过 3 层,我建议在写入 Parquet 前先做一层中间宽表,把 JSON 的 path 直接作为列名。这个操作看似丢失了嵌套信息的简洁性,但在查询时能减少一次列裁剪的递归计算。
- 分配列顺序:Parquet 文件里列的顺序不是随意的。把常被过滤的列放在前面,把大字段(如长文本、图片 URL)放在末尾,读取时扫描头部的统计信息就能提前剔除大块数据。
以下是一个我常用的最终 schema 设计示例,如果你做订单域,可以直接套:
import pyarrow as pa schema = pa.schema([ pa.field("order_id", pa.string()), pa.field("order_date", pa.timestamp("ms")), pa.field("customer_id", pa.int64()), pa.field("customer_name", pa.string()), pa.field("items", pa.list_(pa.struct([ pa.field("product_id", pa.int64()), pa.field("product_name", pa.string()), pa.field("price", pa.decimal128(12, 2)), pa.field("quantity", pa.int32()), ]))), pa.field("total_amount", pa.decimal128(12, 2)), pa.field("tags", pa.list_(pa.string())), pa.field("remark", pa.string()), ])这里有个容易被忽视的点:items这种嵌套结构,Parquet 会拆成多个物理列,但每个嵌套元素都会带 repetition level。如果你完全拍平,可能反而会让文件更小,因为不用存储嵌套结构带来的额外层级信息。我的经验是:嵌套层数不超过 2 层且需要保留对象关系时,保留嵌套;超过 2 层就拍平。
3.2 Parquet 写入与读取的读写分离实践
在我负责的数据平台中,写入链路和读取链路是严格分开的。写入链路通常用 Spark 或 Flink 完成批式落地,读取链路专门用 Trino 或 DuckDB 做查询。你在自己的机器上实验时,可以用 pyarrow 同时承担读写角色,但更好的做法是区分功能模块。
写入侧,我一般这样做:
- 从上游拉取原始 JSON,用 pyarrow 的
Table.from_batches批量累积内存数据。 - 分批写入同一文件的多个 row group。不要等全部数据攒齐才写,那样内存扛不住;可以按 10 万条一个 RecordBatch 持续写。
- 写完文件后,校验文件完整性。最直接的方式是用
pq.read_metadata读取文件底部 footer,确认 schema 和 row group 数量正确。
读取侧的核心能力是谓词下推(predicate pushdown)和列裁剪。你如果基于 pyarrow 写查询,可以这样利用:
import pyarrow.parquet as pq # 只读取必要列 filters = [ ("order_date", ">=", "2024-01-01"), ("customer_id", "in", [1001, 1002, 1003]), ] table = pq.read_table( "orders.parquet", columns=["order_id", "order_date", "customer_id", "total_amount"], filters=filters, ) print(table.to_pandas())这里filters参数会被 pyarrow 翻译成下推到 Parquet 文件级别的过滤逻辑,它会在 row group 的统计信息层先过滤一次,再在列数据层过滤一次。我实际测试下来,过滤效果最强的字段是排序好的低基数字段,尤其是日期字段和枚举字段。
3.3 从“他山之石”到落地:我用 Trino 查 Parquet 的体验
如果你要体会 MPP 数据库查 Parquet 的痛快感,我建议你直接把 Trino 部署起来,或者用 DuckDB 在本地直接跑 SQL。DuckDB 的安装成本低到离谱,一条命令就能启动,而且它能直接查 Parquet 文件,连导入步骤都省了:
SELECT customer_name, sum(total_amount) AS revenue FROM 'orders.parquet' WHERE order_date >= DATE '2024-01-01' GROUP BY customer_name ORDER BY revenue DESC LIMIT 10;你注意,DuckDB 不需要先把数据加载进表里,它直接在文件层面做查询。这是“他山之石”给今天的最大启示:数据存储形态和查询引擎彻底解耦,Parquet 成了标准的静态数据交换格式。MPP 数据库当年强调的分布式并行执行、向量化计算、列裁剪,如今在 Parquet 这张“冷文件”上都能复现。区别仅仅是:MPP 把数据常驻集群内存,而 Parquet 查询是边读边算,靠列存压缩和统计信息把磁盘 IO 降到最低。
我实际用 DuckDB 查过一个 1.2GB 的 Parquet 文件,查询只扫 4 列,耗时从直接读 CSV 的 15 秒降到了 0.8 秒。这里面没有预聚合、没有专门索引,纯粹就是列式文件格式 + 向量化执行引擎带来的收益。这个体验比直接在 Python 里遍历 DataFrame 要爽得多。
4. 常见问题与排查技巧实录
4.1 明明有数据,为什么读出来空结果?
这是我接手过最多的一个坑。很多时候问题出在过滤条件写得太随意,尤其是日期和字符串类型不匹配。比如你明明知道order_date是 timestamp 类型,却用字符串"2024-01-01"去过滤。pyarrow 的过滤 API 对类型非常严格,字符串和 timestamp 不能直接比对。你需要先做一次类型转换:
import pyarrow.compute as pc date_filter = pa.scalar("2024-01-01", type=pa.timestamp("ms")) table = pq.read_table( "orders.parquet", filters=[("order_date", ">=", date_filter)], )如果你一开始创建的 schema 是timestamp("ms"),那过滤值就必须是timestamp("ms")类型。很多时候同一个文件,用 Trino 查是有数据的,但用 pyarrow 直接过滤就返回空,就是因为类型映射不一致。建议你在读取文件后先打印table.schema,确认物理类型和逻辑类型,再写过滤条件。
4.2 文件很大但查询依旧慢,问题出在哪?
如果你发现 Parquet 文件没有发挥出列存的优势,先从这三个方向排查:
- 列裁剪是否生效:用
EXPLAIN看执行计划,确认只扫了需要的列,而不是全列扫描。比如 DuckDB 里用EXPLAIN SELECT ...能看到扫描算子读取的列清单。 - 谓词下推是否生效:检查过滤字段是否出现在 row group 的统计信息里。如果过滤字段是压碎的字符串字段,统计信息可能退化成全量扫描。此时要么对该字段做排序,要么把过滤字段单独建立索引列。
- row group 是否过大或过小:row group 太小(比如默认 128MB)会导致元数据占比升高,查询时每个文件要读大量 footer 元数据;row group 太大(比如 1GB)则会导致单次顺序读过多无关数据。我实测 512MB 是大多数场景下的甜点。
4.3 schema 变更后读不出旧文件怎么办?
工业级数据链路里,schema 变更是家常便饭。Parquet 支持向后兼容,但你在实际运维中会发现:新增字段没问题,删除字段也没有问题,但字段类型变更(比如 string 变成 int64)会引发他们所说的“schema evolution”问题。我的经验是:统一走“先加新列,再迁移旧值”的策略,不要直接改旧字段类型。可以用 SQL 做一个轻量迁移:
CREATE TABLE new_orders AS SELECT order_id, CAST(order_date AS VARCHAR) AS order_date_str, customer_name, total_amount FROM 'orders.parquet';再把这个新查询结果写回 Parquet。如果数据量太大,就用 Spark 的重分区处理,别在单机内存里硬撑。
4.4 嵌套数据写进去之后,为什么查询时 null 变多了?
这个坑很隐蔽,尤其在你用 JSON 转 Parquet 时容易发生。原因在于 JSON 里缺失某个嵌套对象的字段,PyArrow 会补一个null,但如果你没有在 schema 里声明该字段是 optional,写入后读取时会自动补 null。这不是 bug,却往往让人误以为是数据丢失。
我的解决办法是:在数据接入阶段先用pc.fill_null统一填充默认值,避免在查询时才处理。尤其是数值字段,默认值尽量用 0 而不是 null;字符串字段默认值用空串而不是 null。这样的好处是 Parquet 的统计信息更准确,谓词下推时不会因为在某个 row group 里大量 null 导致过滤失效。
5. 从 Dremel 到 Parquet:我的一点个人体会
说实话,刚开始研究 Dremel 和 Parquet 的关系时,我最大的困惑是:MPP 数据库毕竟是一个完整系统,有集群调度、查询优化器、执行引擎,而 Parquet 只是一个文件格式,两者怎么能相提并论?后来我一步步用 Trino、DuckDB 去查询 Parquet,才逐渐意识到,Parquet 是把 MPP 数据库里最核心的“列式存储 + 统计信息 + 数据裁剪”思想抽出来,做成一个标准文件格式,并且让任何查询引擎都能直接消费它。
这相当于把“山”变成了“石”——MPP 数据库是完整山体,Parquet 则是山体上最坚硬、最可复用的石材。任何团队,哪怕只有一台机器,只要把数据落成 Parquet,都能借助 DuckDB、DataFusion 这类轻量引擎快速获得接近 MPP 的查询性能。这一点,在我看来是大数据技术下沉的重要一步。
我个人的建议是:如果你正在设计新的数据仓库或数据湖,不管底层用 Hive Metastore 还是 Iceberg,先把 Parquet 的 schema 设计和 row group 调优做扎实。别一上来就上一套重引擎,先用最小的方案把列存格式的红利吃透。我自己踩过的坑也在这里分享过,比如压缩选型不对、嵌套层级过深、过滤类型不一致,这些都是可以提前规避的。你把这些细节处理好,后面向上扩展也顺滑得多。