Python技术迷

Parquet文件格式,已成为大数据存储的默认选择

那天我在公司楼下买咖啡,排队的时候同事在群里吐槽:“又一个任务跑炸了,S3 上全是 CSV,小文件两百多万个,Spark 读起来跟拖拉机一样”。我当时就笑,心里想:哎这不就是典型的……你还在用“文本文件思维”存大数据嘛。后来回工位我顺手把他那批 CSV 转了 Parquet,第二天他就来一句:卧槽同一段 SQL 怎么快这么多。

你们知道吧,Parquet 现在基本就是大数据存储的默认答案了,不是因为它“潮”,是因为它把一堆真实的坑都给你填上了:

  • 你查报表,多半只用到几列,Parquet 是列式存储,读的时候就真只读那几列(CSV/JSON 你得把一整行整行扫一遍)。
  • 还有压缩,列式天然更好压缩,同一列值类型接近、重复高,压缩比很夸张。
  • 以及最关键的一个:它会带统计信息(min/max/null_count 这种),查询引擎能做跳读,再配合 predicate pushdown(谓词下推),你 where 一加,能少读很多 row group。

我记得我们那次事故是“按天落盘 + 全量扫描”,查最近 7 天订单,却每次都把 3 个月的文件扫一遍。后来我让他先把目录按 dt 分区,再把数据写成 Parquet,基本就安静了。就这么个意思,你别管啥高大上的名词,核心就是:少读、压缩、能跳过、还能并行。

来,给你们一段我自己平时会写的 Python,小而脏但真能跑(你本地得装 pyarrow 和 pandas,别问我为啥你机器没装,问就是环境不干净…)

import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq
from pathlib import Path

defwrite_parquet_demo(out_dir: str):
    out = Path(out_dir)
    out.mkdir(parents=True, exist_ok=True)

# 模拟一点订单数据(别纠结字段,意思到了就行)
    df = pd.DataFrame({
"order_id": range(1, 20001),
"user_id": [x % 3000for x in range(1, 20001)],
"amount":  [float(x % 500) / 10for x in range(1, 20001)],
"dt":      ["2026-02-01"if x % 3 == 0else"2026-02-02"for x in range(1, 20001)],
"status":  ["PAID"if x % 5else"REFUND"for x in range(1, 20001)],
    })

    table = pa.Table.from_pandas(df, preserve_index=False)

# 重点:row_group_size 影响后面“跳读”的颗粒度,别写太小也别太大
    pq.write_table(
        table,
        out / "orders.parquet",
        compression="zstd",        # zstd 很常见,压缩比和速度都不错
        row_group_size=50_000      # demo 数据小,看个意思
    )

write_parquet_demo("./lake_demo")

上面这段写出来的 Parquet,里面就会有 row group、列统计信息这些东西。然后读的时候,你就能“只取列 + 带过滤条件”,这里才是爽点:

import pyarrow.parquet as pq

defread_with_pushdown(path: str):
    pf = pq.ParquetFile(path)

# 只读两列:order_id 和 amount
# 再加过滤:status == 'PAID' 且 amount > 20
# 注意:不同引擎/版本对下推支持程度略有差异,但 Parquet 天生就给你这个能力
    table = pq.read_table(
        path,
        columns=["order_id", "amount", "status"],
        filters=[("status", "==", "PAID"), ("amount", ">", 20.0)]
    )

# 转 pandas 看看
    df = table.to_pandas()
return df

df2 = read_with_pushdown("./lake_demo/orders.parquet")
print(df2.head())
print("rows:", len(df2))

你看这个思路,跟 CSV 完全不一样。CSV 你过滤得先读进来再过滤;Parquet 是尽量让存储层+执行引擎提前帮你“别读不该读的东西”。所以在 Spark、Hive、Trino/Presto 这些体系里,它就特别舒服:

  • 读列少,IO 少;
  • 压缩省磁盘,也省网络(对象存储上省得更明显);
  • row group + 统计信息让引擎能跳过一坨数据;
  • 再加上它支持 schema(字段类型、嵌套结构),“这列到底是 int 还是 string”这种扯皮少很多。

哦对,还有一个我经常在群里骂人的点:小文件。Parquet 再好,小文件照样把你搞死(对象存储 list/seek 成本高,元数据压力也大)。很多人一边说自己用 Parquet,一边每天落几十万文件,我就……你这不是 Parquet 的锅,是你写入策略的问题。

给个我自己“合并小文件”的土办法示例,思路是把目录里多个 Parquet 合成更大的 row group/文件(生产里你可能用 Spark compact 或者 Iceberg/Delta 的 optimize,这里先用纯 Python 表达一下):

from pathlib import Path
import pyarrow as pa
import pyarrow.dataset as ds
import pyarrow.parquet as pq

defcompact_parquet(input_dir: str, output_file: str, target_rows: int = 500_000):
    dataset = ds.dataset(input_dir, format="parquet")
    scanner = dataset.scan()  # 读全量(演示用,真实要按分区处理)

    batches = []
    total = 0

for batch in scanner.to_batches():
        batches.append(batch)
        total += batch.num_rows
if total >= target_rows:
            table = pa.Table.from_batches(batches)
            pq.write_table(table, output_file, compression="zstd", row_group_size=128_000)
            batches, total = [], 0
break

# 如果没到 target_rows,也写一把
if batches:
        table = pa.Table.from_batches(batches)
        pq.write_table(table, output_file, compression="zstd", row_group_size=128_000)

# compact_parquet("./some_parquet_dir", "./compacted.parquet")

反正你就记住:Parquet 不是“把 CSV 改个后缀”,它更像是“数据仓库的地基砖”。你要配合分区策略(比如 dt=2026-02-02 这种目录分区)、文件大小控制(别太碎)、再加上引擎读的时候列裁剪+下推,那个收益才会全出来。

我现在判断一个团队数据链路成熟不成熟,有时候都不看他用啥框架,就看两点: 1)落盘是不是列式(Parquet/ORC 这种) 2)有没有在意分区和文件大小 只要这俩是对的,后面你用 Spark 还是 Trino,甚至你迁湖仓一体,基本都不会太痛。

行了我先不扯了,刚刚看群里又有人说“Parquet 读出来怎么有些列变成了 object”,八成是 schema 演进没管好……我去回两句,不然他今晚又要通宵。