1. 为什么你需要PyArrow和Parquet?

如果你正在用Python处理数据,尤其是数据量稍微大一点,比如超过百万行,或者数据列特别多,你肯定对Pandas又爱又恨。爱的是它语法简洁,功能强大;恨的是它动不动就内存爆炸,处理大文件慢得让人想砸键盘。我以前处理一个几GB的CSV文件,光是pd.read_csv就得等上好几分钟,内存占用直接飙到十几个G,机器卡得不行。

后来我接触到了Parquet格式,配合PyArrow这个库,简直像是打开了新世界的大门。处理同样的数据,速度能快上几倍甚至几十倍,内存占用也大幅下降。这可不是什么“未来科技”,而是现在很多大数据框架(比如Spark、Dask)都在用的标准列式存储格式。

简单来说,Parquet文件就像是一个设计精良的衣柜。传统的CSV文件(行式存储)好比是把所有衣服(数据)一件接一件平铺在箱子里,想找一件衬衫,得把所有衣服都翻一遍。而Parquet(列式存储)则是把衬衫、裤子、袜子分别放在不同的格子里。当你只需要分析“衬衫的颜色”这个信息时,Parquet可以直接打开“衬衫”这个格子,只读取颜色那一列的数据,完全不用管裤子袜子。这种特性,让它在大数据查询和分析场景下效率极高。

PyArrow就是Python里用来高效读写这个“智能衣柜”的超级工具。它底层用C++实现,速度飞快,并且和Pandas是无缝集成的。你完全可以用你熟悉的Pandas语法来操作,但背后却是PyArrow引擎在提供火箭般的加速。接下来,我就带你从最基础的安装开始,一步步掌握如何用PyArrow高效玩转Parquet文件,直到能处理海量数据集。

2. 环境搭建与基础读写操作

2.1 快速安装与配置

万事开头难,但安装PyArrow一点也不难。我强烈推荐使用Conda来管理环境,它能很好地处理一些底层依赖。打开你的终端(或Anaconda Prompt),跟着我敲命令就行:

# 创建一个新的虚拟环境,专门用于数据科学项目
conda create -n data_processing python=3.11
# 激活这个环境
conda activate data_processing
# 安装PyArrow和Pandas
conda install pyarrow pandas

如果你习惯用pip,也可以:

pip install pyarrow pandas

安装过程通常很顺利。完成后,在Python脚本或Jupyter Notebook里导入它们,就可以开始了:

import pyarrow as pa
import pyarrow.parquet as pq
import pandas as pd
print(f"PyArrow版本: {pa.__version__}, Pandas版本: {pd.__version__}")

2.2 你的第一个Parquet文件:写入与读取

让我们从一个最简单的例子开始。假设你有一个Pandas的DataFrame,这是数据分析的起点。

# 创建一个示例DataFrame
df = pd.DataFrame({
    'user_id': [1001, 1002, 1003, 1004],
    'user_name': ['Alice', 'Bob', 'Charlie', 'Diana'],
    'age': [25, 30, 35, 28],
    'score': [88.5, 92.0, 79.5, 95.5],
    'city': ['Beijing', 'Shanghai', 'Guangzhou', 'Shenzhen']
})
print("原始DataFrame:")
print(df)

现在,我们要把这个DataFrame保存为Parquet文件。用PyArrow来做,只需要两步:

# 第一步:将Pandas DataFrame转换为PyArrow的Table格式。
# Table是PyArrow的核心数据结构,类似于内存中的关系表。
table = pa.Table.from_pandas(df)

# 第二步:将Table写入Parquet文件。
# `compression='snappy'` 是默认的压缩方式,在速度和压缩比之间取得了很好的平衡。
pq.write_table(table, 'users.parquet', compression='snappy')

print("Parquet文件 'users.parquet' 已保存。")

写入之后,我们再来读取它。读取的方式非常直观:

# 方法一:直接读取整个文件到Pandas DataFrame(适合小文件)
df_from_parquet = pd.read_parquet('users.parquet', engine='pyarrow')
print("\n通过pd.read_parquet读取:")
print(df_from_parquet)

# 方法二:使用PyArrow的ParquetFile对象(更灵活,适合后续高级操作)
parquet_file = pq.ParquetFile('users.parquet')
table_from_file = parquet_file.read()
df_from_table = table_from_file.to_pandas()
print("\n通过pq.ParquetFile读取:")
print(df_from_table)

这两种方法对于小文件来说结果一样。但第二种方法pq.ParquetFile给了我们一个文件句柄,后面做迭代读取、查看元数据等操作时会用到。你可以用系统命令行看看生成的文件大小,通常比同数据的CSV文件小很多,这就是列式压缩的威力。

3. 进阶数据操作与性能初探

3.1 不只是读取:筛选、转换与导出

读取数据只是第一步,我们通常需要做过滤和计算。PyArrow配合Pandas,能让你在享受高性能的同时,使用熟悉的语法。

假设我们想找出分数高于90分的用户,并计算他们的“年龄分数比”(一个虚构的指标):

# 读取数据
table = pq.read_table('users.parquet')
df = table.to_pandas()

# 使用Pandas语法进行筛选和计算
high_score_df = df[df['score'] > 90].copy()
high_score_df['age_score_ratio'] = high_score_df['age'] / high_score_df['score']

print("高分用户信息:")
print(high_score_df)

这里有个细节:我用了.copy()来创建筛选后数据的一个副本。这是一个好习惯,可以避免后续操作中可能出现的SettingWithCopyWarning警告。

处理完的数据,你可能需要导出为其他格式,比如CSV,用于和只用Excel的同事协作。这也很简单:

# 将处理后的DataFrame保存为CSV
high_score_df.to_csv('high_score_users.csv', index=False)
print("数据已导出为 'high_score_users.csv'")

# 当然,你也可以将处理后的数据保存为新的Parquet文件,保持高效存储
new_table = pa.Table.from_pandas(high_score_df)
pq.write_table(new_table, 'high_score_users.parquet')
print("数据也已保存为 'high_score_users.parquet'")

3.2 处理复杂嵌套数据:从“一坨”到“一列”

在实际项目中,你可能会遇到一些“不规整”的数据。比如,原始数据里有一个feature列,里面存储的不是单个值,而是一个Python列表(List)。这在机器学习的数据集中很常见,比如一个样本的特征向量。直接打印出来,你会看到类似[0.1, 0.5, 0.3, ...]挤在一个单元格里,分析起来非常麻烦。

我遇到过这种情况,一个Parquet文件的feature列有784个浮点数,代表一张28x28图片的像素。我需要把它们拆成784个单独的列。用PyArrow和Pandas可以这样优雅地解决:

# 假设我们读取了一个包含'feature'列表列的文件
# 这里我们模拟创建这样一个DataFrame
import numpy as np
data = {
    'id': [1, 2],
    'feature': [np.random.rand(5).tolist(), np.random.rand(5).tolist()] # 模拟5维特征
}
df_complex = pd.DataFrame(data)
print("原始带列表列的数据:")
print(df_complex)
print(f"feature列第一个元素: {df_complex['feature'].iloc[0]}")

# 关键操作:将列表列拆分成多个单独的特征列
# 使用apply将每个列表转换成一个Series,然后concat合并回原DataFrame
split_features = df_complex['feature'].apply(pd.Series)
# 为拆分出的列命名,例如 feature_0, feature_1...
split_features.columns = [f'feature_{i}' for i in range(split_features.shape[1])]

df_expanded = pd.concat([df_complex.drop('feature', axis=1), split_features], axis=1)
print("\n拆分特征列后的数据:")
print(df_expanded)

这个apply(pd.Series)技巧非常实用,它能将一列列表“爆炸”成多列。之后你就可以方便地对每一个特征维度进行统计分析了。处理完记得把原始的feature列删掉,节省空间。

4. 应对大数据挑战:迭代读取与内存优化

4.1 为什么需要迭代读取?

当你兴冲冲地用pd.read_parquet去打开一个10GB的Parquet文件时,大概率会收到一个“MemoryError”(内存错误)。这是因为Pandas默认会将整个文件加载到内存中,形成DataFrame。对于远超物理内存的大文件,这显然是行不通的。

这时候,就需要用到迭代读取(分批读取)。PyArrow的iter_batches方法就是为此而生的。它允许你像“吃自助餐”一样,一次只拿一盘(一个批次),吃完(处理完)再拿下一盘,而不是把整个餐厅的食物都堆到桌子上。

4.2 实战:分批处理超大Parquet文件

我们来模拟处理一个大文件。假设我们有一个large_data.parquet文件,行数很多。我们分批读取,并对每一批数据进行处理(比如我们之前做的特征列拆分),最后将结果汇总或保存。

import pyarrow.parquet as pq
import pandas as pd
import time

# 记录开始时间
start_time = time.time()

# 1. 创建ParquetFile对象,它代表文件但不立即加载数据
parquet_file = pq.ParquetFile('large_data.parquet') # 请替换为你的大文件路径

# 2. 设置批次大小。这是性能调优的关键参数!
# 太小会导致I/O次数过多,太大会导致单批内存占用过高。
# 根据你的数据列数和内存大小,从1万到10万行开始尝试。
batch_size = 50000

# 3. 获取迭代器
batch_iterator = parquet_file.iter_batches(batch_size=batch_size)

# 4. 初始化一个列表,用于收集每批处理后的DataFrame
processed_dfs = []

# 5. 开始迭代处理
for i, batch in enumerate(batch_iterator):
    # 将当前批次(RecordBatch)转换为Pandas DataFrame
    df_batch = batch.to_pandas()
    
    # 在这里执行你的数据处理逻辑,例如拆分feature列
    # 假设df_batch有'feature'列需要拆分
    if 'feature' in df_batch.columns:
        split_features = df_batch['feature'].apply(pd.Series)
        df_batch = pd.concat([df_batch.drop('feature', axis=1), split_features], axis=1)
    
    # 将处理好的批次DataFrame加入列表
    processed_dfs.append(df_batch)
    
    # 可选:每处理一定批次,打印进度,让自己安心
    if (i + 1) % 10 == 0:
        print(f"已处理 {i + 1} 个批次...")

# 6. 将所有批次的DataFrame合并成一个大的DataFrame
# 注意:如果最终结果仍然很大,这一步也可能内存不足。可以考虑分批写入磁盘。
final_df = pd.concat(processed_dfs, ignore_index=True)

end_time = time.time()
print(f"\n数据处理完成!总耗时: {end_time - start_time:.2f} 秒")
print(f"总数据行数: {len(final_df)}")

通过这种方式,即使原始文件有100GB,只要每个批次的大小在你的内存承受范围内(比如2GB),你就能顺利处理。这是处理海量数据最基本也是最重要的技能之一。

5. 高效处理多文件与生产级优化

5.1 合并多个Parquet文件

在实际生产环境中,数据通常不是放在一个巨大的文件里,而是按日期、按类别分割成成千上万个小的Parquet文件,例如part-00001.parquet, part-00002.parquet... 我们需要将它们合并分析。

最直接的想法是用一个循环,逐个读取再合并。我们来写一个健壮的版本:

import os
import glob
import pyarrow.parquet as pq
import pandas as pd

# 指定包含多个parquet文件的文件夹路径
folder_path = './daily_logs/'  # 假设这里存放了按天划分的日志文件

# 方法1:使用glob匹配所有.parquet文件
parquet_files = glob.glob(os.path.join(folder_path, '*.parquet'))
print(f"找到 {len(parquet_files)} 个Parquet文件")

# 初始化一个空列表,用于存放每个文件读取后的Table
tables = []

for file_path in parquet_files:
    # 读取单个文件为PyArrow Table
    table = pq.read_table(file_path)
    tables.append(table)
    print(f"已加载: {os.path.basename(file_path)}")

# 使用pyarrow.concat_tables将所有Table在行方向合并
# 这比在Pandas层面合并(pd.concat)更高效,因为避免了中间转换。
if tables:
    combined_table = pa.concat_tables(tables)
    combined_df = combined_table.to_pandas()
    print(f"\n合并完成。总数据行数: {len(combined_df)}")
else:
    print("未找到Parquet文件。")

这里我用了pa.concat_tables,它是在PyArrow的层面进行合并,效率高于先转成多个Pandas DataFrame再用pd.concat。因为pd.concat会涉及大量的数据拷贝。

5.2 性能调优核心技巧

当你开始处理真正的大数据时,一些细微的参数调整会带来巨大的性能差异。下面是我踩过坑后总结的几个关键点:

1. 选择正确的压缩算法: Parquet支持多种压缩格式。在pq.write_table时可以通过compression参数指定。

  • snappy (默认):压缩和解压速度非常快,压缩率适中。绝大多数场景的首选,平衡了速度和空间。
  • gzip:压缩率更高,能生成更小的文件,但压缩和解压速度比snappy慢。适合需要长期归档、对磁盘空间敏感,且不常读取的场景。
  • zstd:较新的算法,旨在提供比zlib/gzip更高的压缩比和接近snappy的速度。越来越受欢迎。
  • none:不压缩。除非你后续有特殊的流式处理需求,否则不建议使用。
# 写入时指定压缩算法
pq.write_table(table, 'data_compressed.parquet', compression='zstd')

2. 调整batch_size(批次大小): 前面迭代读取时提到的batch_size是性能关键。它决定了每次从磁盘加载到内存的数据量。

  • 设置过小:比如1000行。会导致函数调用和I/O次数过多,增加总耗时。
  • 设置过大:比如500万行。单批数据可能超过内存,导致程序崩溃。
  • 如何设定:没有银弹。你需要根据数据列的多少(宽度)和可用内存来测试。一个经验法则是,让单批数据在内存中的大小约为可用内存的1/10到1/5。你可以先读取文件的一小部分,估算单行大小,然后计算合适的行数。

3. 使用columns参数进行列裁剪: 这是列式存储最大的优势!如果你只需要用到100列中的5列,那么在读取时就指定它们,可以跳过其他95列的磁盘I/O和解压计算,极大提升速度。

# 只读取‘user_id’, ‘age’, ‘score’三列
df_filtered_columns = pd.read_parquet(
    'large_data.parquet',
    engine='pyarrow',
    columns=['user_id', 'age', 'score'] # 关键参数!
)
print(df_filtered_columns.head())

4. 使用filters参数进行谓词下推: 更厉害的是,你可以在读取时指定过滤条件。PyArrow会尽可能在读取数据时应用这些过滤,提前过滤掉不需要的行,减少传输到内存的数据量。

# 只读取‘city’为‘Shanghai’且‘score’大于90的记录
# 注意:过滤条件的语法是列表嵌套元组
df_filtered_rows = pd.read_parquet(
    'large_data.parquet',
    engine='pyarrow',
    filters=[
        ('city', '=', 'Shanghai'),
        ('score', '>', 90)
    ]
)
print(f"筛选后记录数: {len(df_filtered_rows)}")

谓词下推和列裁剪是生产环境中优化查询性能最有效的手段,务必掌握。

6. 避坑指南与最佳实践

在我多年的使用中,积累了一些“血泪教训”,希望能帮你少走弯路。

坑1:数据类型映射问题 PyArrow和Pandas的数据类型并非完全一一对应。比如,PyArrow的timestamp精度可以到纳秒,而Pandas的datetime默认是微秒。在复杂数据类型(如嵌套结构)来回转换时,可能会丢失信息或报错。

  • 建议:在转换后,检查一下df.dtypes,确保关键列的类型符合预期。对于时间戳,可以使用pd.to_datetime进行强制转换。

坑2:字符串列带来的内存膨胀 Pandas中,字符串列(object dtype)的内存效率并不高。如果Parquet文件中有很长的文本字段,用Pandas读取后内存占用可能会激增。

  • 建议:如果不需要对这些字符串列进行计算,在读取时用columns参数排除它们。或者,考虑在PyArrow Table层面进行字符串过滤操作,只在最后必要步骤转为Pandas。

坑3:盲目追求极致压缩 使用gzipzstd的高压缩级别可以得到更小的文件,但写入时间会成倍增加。对于需要频繁写入的中间数据,这很不划算。

  • 建议:对最终归档的、不常变动的数据使用高压缩率算法。对ETL管道中的中间临时文件,使用snappy甚至不压缩。

坑4:忽略分区数据集 对于海量数据,最高效的方式是使用分区。你可以按照日期(year=2023/month=10/day=05)、地区等维度将数据存储在不同的子文件夹中。PyArrow可以高效地读取分区数据集,并且结合filters参数,可以做到“按需读取”,性能提升是指数级的。

  • 建议:如果你的数据量达到TB级别,一定要研究并使用分区功能。pyarrow.dataset模块提供了强大的分区数据集读写接口。
# 一个分区数据集的读取示例(假设数据按‘date’和‘region’分区)
import pyarrow.dataset as ds

dataset = ds.dataset('s3://my-bucket/logs/', format='parquet') # 也支持本地路径
# 只读取2023年10月,且区域为‘east’的数据
condition = (ds.field('date') >= '2023-10-01') & (ds.field('date') < '2023-11-01') & (ds.field('region') == 'east')
table = dataset.to_table(filter=condition)

掌握这些基础、进阶和优化的知识后,你基本上就能应对日常工作中绝大多数与Parquet文件打交道的情景了。从快速处理几GB的数据,到搭建能处理TB级数据的稳健管道,PyArrow都是你手中不可或缺的利器。记住,关键是多实践,根据自己数据的特性去调整参数,找到最适合你那个场景的“甜点”配置。

更多推荐