You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

提升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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 12:20:31