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

为何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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 05:43:04