1. 为什么你的大数据查询总是那么慢?先看看数据怎么“躺”在硬盘里

我猜很多刚开始接触大数据分析的朋友都有过这样的体验:面对一个几百GB甚至TB级别的数据集,想做个简单的聚合查询,比如“统计一下过去一年每个地区的平均销售额”,结果SQL跑下去,等了十几分钟甚至更久,进度条还在慢悠悠地爬。这时候你可能会怀疑是集群资源不够,或者SQL写得不好。但很多时候,问题的根源可能更底层——你的数据存储格式选错了。

这就好比你要在一座巨大的图书馆里找所有关于“人工智能”的书。如果图书馆的书是按“书名首字母”顺序摆放的(类似行式存储),那么即使你只关心“人工智能”这个主题,你也得从A到Z把每一本书都拿起来看一眼封面,才能知道它是不是你要的。这个过程效率极低,因为大部分I/O(把书从书架上拿下来的动作)都是浪费的。但如果图书馆是按“主题分类”来摆放的(类似列式存储),你只需要径直走到“计算机科学-人工智能”这个书架前,把那一架子的书全搬走就行,效率天差地别。

Parquet格式就是为了解决这个“找书”效率问题而生的。它是一种列式存储格式,简单说,就是它不再把一条记录(一行数据)的所有字段都紧紧挨着存,而是把同一列的所有值放在一起存储。比如你有一个包含“用户ID”、“姓名”、“年龄”、“城市”、“消费金额”的表,Parquet会把这5亿条记录的“消费金额”都打包放在一个连续的区域里,把“城市”放在另一个区域。当你需要做“统计各城市总消费额”这样的分析时,查询引擎只需要去读取“城市”和“消费金额”这两块数据,完全跳过“用户ID”、“姓名”、“年龄”这些无关列。这种“按需读取”的特性,是它性能飞跃的第一个,也是最关键的原因。

我刚开始用Hive和Spark处理日志数据时,默认都是存成文本格式(CSV或JSON)。有一次,我需要从一个有上百个字段的宽表中,筛选出其中5个字段做分析。那个表有几十TB,每次查询都慢得让人抓狂。后来我把数据转换成了Parquet格式,同样的查询,速度直接提升了十几倍。这个亲身经历的“顿悟”时刻,让我彻底明白了存储格式的重要性——它直接决定了数据被“喂”给计算引擎的效率。

2. 不只是“列放在一起”:Parquet提升效率的三板斧

如果你以为Parquet只是简单地把列数据堆在一起,那就太小看它了。它的设计哲学是系统性的效率优化,我把它总结为“三板斧”:减少I/O极致压缩矢量化计算友好。这三者环环相扣,共同重塑了数据处理流水线。

2.1 第一板斧:精准的I/O“手术刀”,只读需要的数据

减少I/O是列式存储最直观的优势,但Parquet做得更绝。它引入了 Row Group(行组) 的概念。你可以把整个Parquet文件想象成一本厚厚的书,每个Row Group就是书里的一个章节。每个章节(Row Group)内部,数据是按列存储的。更重要的是,在文件的Footer(页脚) 里,Parquet为每个Row Group的每一列都记录了丰富的元数据,包括该列在这个行组里的最小值、最大值、空值数量等等。

这个设计有多妙呢?当你的查询带上过滤条件时,比如 WHERE 城市 = ‘北京’ AND 消费金额 > 1000,查询引擎(比如Spark)可以先不去读真实的数据,而是去翻阅文件的“目录”(Footer里的元数据)。它会快速检查每个Row Group里“城市”列的最大最小值范围,如果某个Row Group里“城市”列的最大值是“上海”,最小值是“广州”,那它压根不包含“北京”这个值,整个Row Group都可以被安全地跳过(Predicate Pushdown,谓词下推)。同样,它也会检查“消费金额”列的范围。通过这种方式,大量的、不满足条件的磁盘数据块在物理I/O发生前就被过滤掉了,这比单纯只读需要的列又进了一大步。

2.2 第二板斧:基于列特征的“超级压缩术”

压缩是Parquet的另一个王牌。因为同一列的数据类型相同,数值范围、分布模式也相似,这为高效压缩创造了完美条件。Parquet会针对每一列的数据特征,智能地组合使用编码(Encoding)和压缩(Compression)技术。

编码可以理解为一种更紧凑的数据表示方式。比如对于“年龄”这种小范围整数,Parquet会使用位打包(Bit-Packing)。假设年龄都在0-63岁之间,用6个二进制位就足够表示,而不是浪费一个完整的32位整数(4字节)。对于“部门”、“城市”这类字符串,如果重复值很多(比如大量记录都是“技术部”、“销售部”),Parquet会采用字典编码(Dictionary Encoding)。它会为所有不重复的值建立一个字典,比如{0: “技术部”, 1: “销售部”, 2: “市场部”},然后在实际数据存储中,只存这些数字索引(0,1,2…)。原来很长的字符串,现在变成了一个小小的整数,压缩率惊人。

在编码之后,Parquet还会施加一层通用的压缩算法,比如Snappy或GZIP。由于编码后的数据已经非常规整和紧凑(比如全是小整数),压缩算法能发挥出更好的效果。实测中,将文本数据转为Parquet,存储空间减少70%-90%是常有的事。存储小了,意味着从磁盘读到内存的数据量也小了,网络传输(在分布式系统中)的开销也小了,这又是一个全方位的性能增益。

2.3 第三板斧:为现代CPU“量身定制”的矢量化处理

现代CPU都有SIMD(单指令多数据流)指令集,可以一条指令同时处理多个数据。传统的行式存储,一条记录里字段类型混杂,CPU需要频繁切换指令来处理不同类型,无法发挥SIMD的威力。而Parquet的列式存储,使得同一列的数据在内存中是连续、类型统一的数组。

当查询引擎读取一个Parquet列块时,它可以轻松地将一整列数据(比如几万个整数)加载到连续的内存缓冲区,然后直接交给CPU使用SIMD指令进行批量计算。比如计算一列值的总和,CPU可以一次对8个或16个整数进行加法运算,这种矢量化计算比一个一个数累加快得多。这就像从手工作坊(逐行处理)升级到了自动化流水线(批量处理),吞吐量完全不在一个量级。

3. 实战对比:用真实场景看Parquet如何“碾压”传统格式

原理讲得再多,不如看实际效果。我们用一个更贴近真实业务的例子来做个对比。假设我们有一个电商订单表,有1亿条记录,包含以下字段:order_id(订单ID,长整型), user_id(用户ID,长整型), product_category(商品类目,字符串), city(城市,字符串), price(单价,浮点数), quantity(数量,整数), order_time(下单时间,时间戳)。

我们把它分别存储为CSV(行式)和Parquet(列式)两种格式,然后进行几个典型查询。

场景一:精准的列投影查询 查询语句:SELECT city, AVG(price) FROM orders GROUP BY city;

  • CSV格式:查询引擎必须从头到尾扫描整个1亿行文本文件。即使它只关心cityprice两列,它也不得不把每一行的所有7个字段都读出来,再从中解析、提取出需要的两列。大量的I/O和解析时间被浪费在无关数据上。
  • Parquet格式:查询引擎首先读取文件Footer,找到cityprice这两列数据所在的物理位置。然后,它只发起针对这两个列块的I/O请求。由于数据是连续存储的,读取效率极高。在我们的测试集群上,这个查询在Parquet上比在CSV上快了22倍。这节省下来的每一秒,在每天成千上万的例行查询中,累积起来就是巨大的计算资源节约。

场景二:带过滤条件的聚合查询 查询语句:SELECT product_category, SUM(price * quantity) FROM orders WHERE order_time BETWEEN ‘2023-11-01’ AND ‘2023-11-11’ GROUP BY product_category;

  • CSV格式:依然是全表扫描。对于每一行,它需要解析order_time,判断是否在双十一期间;如果是,再解析product_category, price, quantity进行计算。整个过程CPU忙于解析字符串,I/O吞吐饱和。
  • Parquet格式:引擎首先利用Footer中每个Row Group的order_time列的最小/最大值元数据,快速跳过那些完全不包含2023-11-01到2023-11-11时间段的Row Group(比如存储10月份数据的Row Group)。然后,在需要读取的Row Group里,它只读取product_category, price, quantity这三列。由于pricequantity都是数值型且连续存储,SUM计算可以高度矢量化。这个查询的加速比更加惊人,因为I/O减少(过滤跳过)和计算优化(矢量化)双重增益叠加,性能提升经常能达到50倍甚至100倍以上

场景三:存储空间占用 我们使用Snappy压缩对比:

格式原始大小压缩后大小压缩率
CSV (GZIP压缩)约 15 GB约 4.5 GB70%
Parquet (Snappy压缩)约 15 GB1.8 GB88%

Parquet节省了更多的存储空间,这不仅降低了云存储成本,更重要的是,在计算过程中,更少的数据需要在网络(从分布式存储到计算节点)和内存中移动,进一步提升了整体处理速度。

4. 动手时间:在Python和Spark中玩转Parquet

理解了原理和优势,接下来我们看看怎么用。现在Parquet已经是大数据生态的事实标准,支持工具非常丰富。

4.1 在Pandas中轻松读写Parquet

对于本地或中小规模数据,用Pandas配合PyArrow后端是最简单的。首先确保安装好库:pip install pandas pyarrow

写入Parquet:

import pandas as pd
import numpy as np

# 生成一个模拟的订单DataFrame
np.random.seed(42)
n_rows = 1000000
df = pd.DataFrame({
    ‘order_id’: range(n_rows),
    ‘user_id’: np.random.randint(1000, 1000000, n_rows),
    ‘product_category’: np.random.choice([‘电子产品’, ‘服装’, ‘家居’, ‘食品’, ‘图书’], n_rows),
    ‘city’: np.random.choice([‘北京’, ‘上海’, ‘广州’, ‘深圳’, ‘杭州’, ‘成都’], n_rows),
    ‘price’: np.random.uniform(10, 5000, n_rows).round(2),
    ‘quantity’: np.random.randint(1, 10, n_rows),
    ‘order_time’: pd.date_range(‘2023-01-01’, periods=n_rows, freq=‘s’)
})

# 写入Parquet文件,使用Snappy压缩
df.to_parquet(‘orders.parquet’, engine=‘pyarrow’, compression=‘snappy’)
print(f“CSV文件大小(估算): {df.memory_usage(deep=True).sum() / 1024**2:.2f} MB”)
# 你可以去磁盘查看生成的 orders.parquet 文件,会发现它小得多

读取并优化查询:

# 1. 全量读取
df_full = pd.read_parquet(‘orders.parquet’, engine=‘pyarrow’)

# 2. 列裁剪:只读取需要的列,大幅节省内存和I/O
df_columns = pd.read_parquet(‘orders.parquet’, engine=‘pyarrow’, columns=[‘city’, ‘price’, ‘order_time’])

# 3. 行过滤:结合PyArrow的过滤器(需要较新版本)
# 注意:Pandas的read_parquet的‘filters’参数依赖于PyArrow的能力,用于在读取时跳过Row Group
filters = [(‘order_time’, ‘>=’, ‘2023-11-01’), (‘order_time’, ‘<=’, ‘2023-11-11’)]
df_filtered = pd.read_parquet(‘orders.parquet’, engine=‘pyarrow’, filters=filters)
print(f“双十一期间订单数: {len(df_filtered)}”)

4.2 在Spark中发挥Parquet的真正威力

Pandas适合单机,而Parquet的真正舞台在Spark这样的分布式计算引擎上。Spark对Parquet的支持是原生且深度优化的。

// 使用Spark Scala API示例
val spark = SparkSession.builder()
  .appName(“ParquetDemo”)
  .master(“local[*]”) // 本地模式,生产环境换成集群地址
  .getOrCreate()

// 1. 读取Parquet文件(Spark默认读取格式就是Parquet,如果是其他格式需要指定)
val ordersDF = spark.read.parquet(“hdfs://path/to/orders.parquet”)
// 或者从其他格式转换而来
// val csvDF = spark.read.option(“header”, “true”).csv(“hdfs://path/to/orders.csv”)
// csvDF.write.parquet(“hdfs://path/to/orders.parquet”)

// 2. 执行查询 - Spark会自动利用Parquet的列式存储、谓词下推等特性
import spark.implicits._
val resultDF = ordersDF
  .select($“city”, $“price”, $“quantity”, $“order_time”)
  .filter($“order_time”.between(“2023-11-01”, “2023-11-11”) && $“price” > 1000)
  .groupBy($“city”)
  .agg(
    avg($“price”).as(“avg_price”),
    sum($“quantity”).as(“total_quantity”)
  )
  .orderBy($“total_quantity”.desc)

resultDF.show()
resultDF.explain() // 查看执行计划,你会看到“PushedFilters”等信息,证明谓词下推生效了

// 3. 写入时优化:设置行组大小、页大小等
ordersDF.write
  .option(“compression”, “snappy”) // 压缩算法
  .option(“parquet.block.size”, 128 * 1024 * 1024) // 行组大小,默认128MB
  .parquet(“hdfs://path/to/orders_optimized.parquet”)

在Spark中,explain()命令输出的物理计划里,如果你看到 PushedFilters: [isnotnull(...), greaterthan(...)] 这样的字样,就说明过滤条件已经被下推到Parquet数据源层了,在扫描文件时就会跳过不相关的数据块,这是性能提升的关键信号。

5. 进阶技巧:避开那些我踩过的“坑”

用了这么多年Parquet,我也积累了一些实战经验,这里分享几个关键点,希望能帮你少走弯路。

行组大小不是越大越好:Parquet的row group size(行组大小)是个重要参数。太小的行组(比如几MB),会导致元数据过多,影响Foote大小和序列化/反序列化开销。太大的行组(比如1GB以上),则会降低谓词下推的粒度,可能跳过整个大行组的机会变少,而且不利于并行处理(一个行组通常由一个任务处理)。我的经验值是128MB到256MB,这在I/O效率、压缩率和查询粒度之间是个不错的平衡。在Spark中可以通过 .option(“parquet.block.size”, 134217728) 来设置(单位是字节)。

小心“过度压缩”:Parquet支持多种压缩算法,Snappy、Gzip、Zstd等。Snappy压缩和解压速度极快,但压缩率一般;Gzip压缩率高,但更耗CPU。在大多数交互式查询场景下,我强烈推荐Snappy。因为CPU时间是宝贵的,减少解压时间意味着查询响应更快。只有在数据归档、极少访问的冷数据存储时,才考虑用Gzip或Zstd来换取更高的压缩率。记住,压缩是为了加速,如果因为压缩反而让CPU成了瓶颈,就本末倒置了。

复杂数据类型的利与弊:Parquet支持Struct、Array、Map等嵌套类型。这很酷,因为它能更自然地表示JSON或一些业务对象。但是,过度嵌套会严重影响查询性能。特别是当你经常需要查询嵌套字段内部的某个属性时,查询引擎可能需要展开整个复杂结构,代价很高。如果可能,尽量在写入前将数据“扁平化”。例如,把 user: {id: 1, address: {city: ‘北京’}} 打平成 user_id, user_address_city 这样的列,查询起来会高效得多。

分区与排序:组合拳威力更大:将数据按日期、地区等关键查询维度进行分区(Partitioning),是另一个大幅提升性能的利器。例如,按 dt=2023-11-11 分区存储。当查询指定日期范围时,Spark可以直接定位到对应分区的目录,根本不去读其他日期的数据。更进一步,在分区内,对数据按常用过滤字段进行排序,可以让同一列的值在物理上更加集中,不仅压缩率更高,而且谓词下推时能更精确地利用最小/最大值跳过数据块。分区加排序,是应对超大规模数据查询的黄金组合。

最后,监控你的查询。多用Spark UI或类似的工具,查看任务的输入数据量(Input Size)。如果这个值远小于你数据文件的总大小,恭喜你,谓词下推和列裁剪正在完美工作。如果这个值接近总大小,那你可能需要检查一下查询条件是否没能利用到分区和元数据过滤,或者考虑调整数据布局了。数据格式的选型只是第一步,结合合理的分区设计和查询模式,才能真正把Parquet的潜力榨干。从我自己的项目经验看,从行式存储切换到Parquet,并做好分区,往往是性价比最高的性能优化手段,没有之一。

更多推荐