
2. 核心细节解析与实操要点2.1 嵌套数据的拍平逻辑与控制参数在动手写代码之前先把我理解的 Dremel 思路讲透否则你会在字段展开和 null 处理上被折磨到怀疑人生。Dremel 的论文里定义了 record 和 column 两种视角核心是把一棵嵌套 JSON 树通过repeated字段展开成多行。对应到 Parquet 物理层就是 record shredding 和 assembly 的过程。你要是不把repeated和required/optional的语义想清楚写出来的 schema 和真实数据会有“对不上”的错位感。我把一个典型订单数据拍平前后做了对比你可以直观感受一下什么叫“列式存储下的嵌套展开”字段路径原始 JSON 形态拍平后的列名重复次数说明order_id普通字符串order_id1每条记录对应一个值items数组元素含 product 和 priceitems.product, items.price随数组长度变化需用 repeated 修饰customer.name嵌套对象customer.name1合并路径即可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 会自动推断成一个listitem: string的嵌套列它在 Parquet 里会被映射成 repeated 字段。如果你用 pandas 直接存很多时候 list 会被解释成 object 列后续查询引擎没法做很好的下推优化。所以我的习惯是先用 pyarrow 构造好原始 table再统一落盘。还有几个参数我在生产中常用一次性给你pq.write_table( table, orders_zstd.parquet, compressionzstd, row_group_size512 * 1024 * 1024, # 512MB version2.6, data_page_version1.0, )这里版本号别乱填。Parquet 2.6 支持更新的逻辑类型比如TIMESTAMP_NANOS如果你只要普通时间戳可以保持默认但如果你要支持微秒或纳秒精度用version2.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], filtersfilters, ) 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, typepa.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 调优做扎实。别一上来就上一套重引擎先用最小的方案把列存格式的红利吃透。我自己踩过的坑也在这里分享过比如压缩选型不对、嵌套层级过深、过滤类型不一致这些都是可以提前规避的。你把这些细节处理好后面向上扩展也顺滑得多。