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

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仪表盘信息

执行代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 18:07:58