为何pandas.read_parquet无法充分利用多CPU?如何提速?
问题解析与优化方案
为什么Pandas/PyArrow单进程读Parquet用不满CPU?
- 单进程并行天花板:PyArrow的
read_table和Pandas的read_parquet(pyarrow引擎)默认是单进程内线程并行,但Parquet读取的部分环节(如PyArrow格式转Pandas对象)受Python GIL限制,无法充分利用96核的多核性能。即便调整io_thread_count,也仅能优化IO操作的线程数,计算转换环节的并行度仍被单进程卡死。 set_io_thread_count未生效的坑:PyArrow的IO线程池在首次初始化后就固定,若在读取操作之后调用pa.set_io_thread_count(),无法修改已创建的线程池配置,必须在任何IO操作前设置线程数才会生效。- 行组粒度不合理:如果Parquet文件的行组过大,PyArrow单进程的线程池无法将任务拆分为足够多的细粒度子任务,导致CPU核心无法被填满。
Dask为啥能高效利用CPU?
Dask采用多进程/多线程分布式架构,天生适配多核场景:
- 自动将Parquet文件按行组分拆为多个分区,每个分区分配给独立的worker进程/线程处理,直接突破单进程的核心数限制。
- 每个分区的读取还会叠加使用PyArrow的IO线程池,形成“多进程+单进程内线程”的双重并行,因此CPU使用率能冲到600%-1000%(此处百分比以单核心100%为基准计算)。
- Dask的任务调度器会自动平衡负载,避免核心闲置。
提升Pandas/PyArrow读取速度的优化方案
1. 正确配置PyArrow并行
务必在读取操作前设置IO线程数,同时使用PyArrow的Dataset API强化并行能力:
import pyarrow as pa import pyarrow.parquet as pq import pandas as pd # 先设置线程数,再执行任何IO操作 pa.set_io_thread_count(32) # 用Dataset API读取,显式开启并行 dataset = pq.ParquetDataset("s3://你的存储桶/目标文件.parquet") table = dataset.read(use_threads=True) df = table.to_pandas()
也可直接给Pandas的read_parquet传递参数:
df = pd.read_parquet("s3://你的存储桶/目标文件.parquet", engine="pyarrow", use_threads=True)
2. 拆分大文件为小文件
单个66GB的Parquet文件粒度太大,拆分为1-5GB的小文件后,PyArrow单进程可同时读取多个文件,提升并行度。用Dask即可快速完成拆分:
import dask.dataframe as dd dask_df = dd.read_parquet("s3://你的存储桶/大文件.parquet", engine="pyarrow") dask_df.to_parquet("s3://你的存储桶/拆分后的文件目录/", engine="pyarrow", write_index=False)
3. 手动用多进程读取行组
通过多进程分别读取不同行组再合并,绕开单进程的并行限制:
import pyarrow.parquet as pq import pandas as pd from concurrent.futures import ProcessPoolExecutor def read_single_row_group(row_group_idx): reader = pq.ParquetFile("s3://你的存储桶/目标文件.parquet") table = reader.read_row_group(row_group_idx) return table.to_pandas() # 获取文件的行组总数 reader = pq.ParquetFile("s3://你的存储桶/目标文件.parquet") total_row_groups = reader.num_row_groups # 多进程读取,worker数可根据核心数调整 with ProcessPoolExecutor(max_workers=16) as executor: df_list = list(executor.map(read_single_row_group, range(total_row_groups))) # 合并所有结果 final_df = pd.concat(df_list, ignore_index=True)
4. 升级PyArrow版本
你当前使用的PyArrow 12.0.0版本较旧,后续版本(如15.x及以上)对Parquet读取的多核并行逻辑做了不少优化,尤其是数据转换和行组处理环节,升级到最新稳定版大概率能提速:
pip install --upgrade pyarrow
5. 优化S3访问配置
确保在同区域访问S3,必要时开启S3传输加速,减少IO延迟对读取速度的影响。
内容的提问来源于stack exchange,提问作者Bewang
相关产品推荐
相关产品推荐

