如何用Dask并行化将1600万行CSV转换为Parquet格式?
优化Dask CSV转Parquet的并行效率
哇,2小时确实有点离谱——Dask本来就是为并行处理设计的,肯定有优化空间!我来给你几个实用的调整方向,应该能把时间砍下来一大截:
1. 先把Dask的并行能力拉满
默认情况下,Dask可能没有充分利用你的CPU核心。试试启动一个本地分布式集群,手动控制工作进程和线程数,让它跑满你的硬件:
from dask.distributed import Client # 根据你的CPU核心数调整,比如8核机器可以设n_workers=4, threads_per_worker=2 client = Client(n_workers=4, threads_per_worker=2)
这样Dask会更高效地调度任务,避免单线程瓶颈。
2. 优化CSV读取环节
读取是第一个大瓶颈,这里有两个关键调整:
- 手动指定分区大小:默认的
read_csv分区可能太大或太小,用blocksize参数控制每个分区的大小(比如100MB),让并行任务数量更合理:
你的2G文件会分成20个左右的分区,刚好适合并行处理。df = ddf.read_csv('data.csv', blocksize='100MB') - 跳过类型推断:Dask默认会逐列推断数据类型,这对1600万行的文件来说非常耗时。提前手动指定每列的
dtype:
这一步能省掉大量的推断时间,效果非常明显。dtype_dict = { # 替换成你的22列实际类型,比如: 'user_id': 'int32', 'order_amount': 'float64', 'order_date': 'object', # ... 其他列 } df = ddf.read_csv('data.csv', blocksize='100MB', dtype=dtype_dict)
3. 优化Parquet写入环节
写入时的参数调整也能大幅提速:
df.to_parquet( 'data.pq', engine='pyarrow', # pyarrow在多数场景下比fastparquet更快 write_index=False, # CSV没有索引,不用写入浪费时间 compression='snappy', # Snappy压缩速度快,兼顾压缩比 row_group_size=1_000_000, # 每个行组100万行,适配Parquet的存储优化 overwrite=True # 如果目标文件存在直接覆盖 )
- 引擎选择:优先用
pyarrow,它的写入性能普遍优于fastparquet; - 压缩算法:Snappy比Gzip快很多,虽然压缩比稍低,但对于2G文件来说,最终Parquet文件大小也不会太大;
- 行组大小:设置成100万行左右,既能保证并行写入的效率,也方便后续读取时的性能。
4. 排查硬件瓶颈
如果以上调整后还是慢,看看是不是磁盘IO拖了后腿——机械硬盘的读写速度远不如SSD,换成SSD的话,IO瓶颈会瞬间缓解很多。
把这些优化点整合起来,你的转换时间应该能从2小时降到几分钟到几十分钟(具体看你的CPU和磁盘配置)。
内容的提问来源于stack exchange,提问作者Daniel Mahler
相关产品推荐
相关产品推荐

