Dask多DataFrame按列合并(axis=1)时性能与资源消耗过高的优化求助
Dask多DataFrame按列合并(axis=1)时性能与资源消耗过高的优化求助
问题背景
我刚接触Dask,尝试对一组DataFrames执行按列合并(concat,axis=1)时,发现耗时、资源占用和任务数量都远超预期。以下是我的运行细节:
- 调度器与客户端同节点
- 2个工作节点
- 有6个按分区存储的Parquet目录,每个目录包含4个分区文件;最终输出是25行、87000列的数据。每个输入文件夹包含1-10行、10000到20000列,单个目录大小都小于1GB。
想请教下我哪里可能出错了,以及如何减少耗时、资源占用和任务步骤?
警告日志
UserWarning: Sending large graph of size 11.14 MiB. This may cause some slowdown. Consider loading the data with Dask directly or using futures or delayed objects to embed the data into the graph without repetition. See also https://docs.dask.org/en/stable/best-practices.html#load-data-with-dask for more information. warnings.warn( client.py:3371: UserWarning: Sending large graph of size 10.33 MiB. This may cause some slowdown. Consider loading the data with Dask directly or using futures or delayed objects to embed the data into the graph without repetition. See also https://docs.dask.org/en/stable/best-practices.html#load-data-with-dask for more information. warnings.warn( Duration at the end 803.7652008533478 seconds completed writing to file
Dask仪表盘信息
- 仪表盘截图:

- 仪表盘视频:Dask仪表盘视频
执行代码
import dask.dataframe as dd from dask.distributed import Client import sys import os import glob import time # Start the timer start_time = time.time() # Connect to the Dask distributed cluster client = Client('IP:8786') # Replace with your scheduler address and port dirs = sys.argv[1] directory = dirs # Use a wildcard pattern to get all Parquet file paths in the directory parquet_files = glob.glob(os.path.join(directory, '*.parquet')) print(parquet_files) df_list = [dd.read_parquet(file) for file in parquet_files] df_list = [df.set_index('ID') for df in df_list] df_list = [df.persist() for df in df_list] concatenated_df = dd.concat(df_list, axis=1) output_path = 'output.parquet' # Write the DataFrame to Parquet files in parallel try: concatenated_df.to_parquet(output_path, write_index=True) except Exception as e: print(f"Error writing Parquet files: {e}") raise # Print duration and close client end_time = time.time() print(f"Duration at the end {end_time-start_time} seconds") print("completed writing to file") client.close()
补充说明
我现在用的是小数据集测试,我知道用Pandas处理会更快更简单,但我的目标是后续扩展到超大型DataFrames,所以先拿小数据做验证。
期望:能更快运行代码,合理使用资源,减少任务执行数量。
备注:内容来源于stack exchange,提问作者sandeysh
相关产品推荐
相关产品推荐

