提升PyArrow读取性能:Parquet分片转Table慢的优化问题
如何加速PyArrow Dataset分片转Table的速度?
问题背景
我在内部S3云上存储了一个分区数据集,使用PyArrow Table读取该数据集,代码如下:
import pyarrow.dataset as ds my_dataset = ds.dataset(ds_name, format="parquet", filesystem=s3file, partitioning="hive") fragments = list(my_dataset.get_fragments()) required_fragment = fragments.pop()
分片元数据
所需分片的元数据如下:
required_fragment.metadata
输出结果:
<pyarrow._parquet.FileMetaData object at 0x00000291798EDF48> created_by: parquet-cpp-arrow version 9.0.0 num_columns: 22 num_rows: 949650 num_row_groups: 29 format_version: 1.0 serialized_size: 68750
性能问题
但将该分片转换为Table耗时极长:
%timeit required_fragment.to_table()
输出结果:
6min 29s ± 1min 15s per loop (mean ± std. dev. of 7 runs, 1 loop each)
Table本身大小约为272MB:
required_fragment.to_table().nbytes
输出结果:
272850898
请问有什么方法可以加快ds.fragment转Table的速度?
更新
我尝试改用pyarrow.parquet而非pyarrow.dataset,仅修改部分代码:
import pyarrow.parquet as pq my_dataset = pq.ParquetDataset(ds_name, filesystem=s3file, partitioning="hive", use_legacy_dataset=False ) fragments = my_dataset.fragments required_fragment = fragments.pop()
再次测试后性能大幅提升:
%timeit required_fragment.to_table()
输出结果:
12.4 s ± 1.56 s per loop (mean ± std. dev. of 7 runs, 1 loop each)
虽然性能提升令人满意,但我感到困惑,因为设置use_legacy_dataset=False后,底层逻辑应与pyarrow.dataset类似。
电脑信息
- 安装内存:21.0GB
- 操作系统:Windows 10 Enterprise
- 网络速度:10Mbps(下载)/156Mbps(上传)
- S3存储位置:亚洲
解决方案及差异分析
1. 底层引擎的实际差异
尽管官方说明use_legacy_dataset=False会复用pyarrow.dataset引擎,但两者在分片处理的细节上仍存在区别:
pyarrow.dataset的get_fragments()默认会触发全量元数据扫描,而pq.ParquetDataset可能更高效地缓存了元数据,避免重复从S3读取元数据文件- 并行读取策略不同:
pq.ParquetDataset默认启用了行组级别的并行读取,而pyarrow.dataset的to_table()默认配置更保守,需手动开启并行
2. pyarrow.dataset的优化方法
- 开启并行读取:在
to_table()中显式指定use_threads=True,利用多线程加速数据读取:required_fragment.to_table(use_threads=True) - 预加载全量元数据:初始化数据集时设置
metadata_mode="full",提前加载所有元数据,避免转换时重复读取:my_dataset = ds.dataset(ds_name, format="parquet", filesystem=s3file, partitioning="hive", metadata_mode="full") - 按需读取列:如果不需要全部22列,仅读取业务所需列,减少数据传输量:
required_fragment.to_table(columns=["所需列1", "所需列2"]) - 优化S3文件系统配置:调整S3连接的超时时间、重试次数等参数,提升云存储访问稳定性:
import pyarrow.fs as fs s3file = fs.S3FileSystem(connect_timeout=30, read_timeout=30, max_attempts=5)
3. 性能差异的补充说明
pq.ParquetDataset作为Parquet格式专用的读取工具,在初始化时会针对Parquet的行组、列存储特性做针对性优化;而pyarrow.dataset是通用数据集引擎,默认配置需要兼顾多种格式,因此手动调整参数后才能达到与pq.ParquetDataset相近的性能。
内容的提问来源于stack exchange,提问作者Femi King
相关产品推荐
相关产品推荐

