Python中Dask Series compute()执行缓慢及比值计算函数验证求助
处理大规模数据列间两两比值的Dask性能问题
需求说明
需要处理大规模数据,按列计算两两比值:即列中第i个值依次除以第i+1及之后的所有值。
示例输入数据
0 1 2 0 34.04 56.55 49.65 1 35.86 49.28 42.36 2 11.21 33.96 18.11 3 12.40 23.17 16.28 4 1087.51 93.37 166.75 5 494.39 182.15 893.23 6 11018.44 6044.63 17347.33 7 20.48 38.01 92.82 8 14866.34 9034.35 19625.89 9 21932.70 12289.75 43752.48 10 8561.40 3279.36 8717.61 11 1050.82 1302.27 1951.25 12 978.03 202.63 179.67 13 15.22 28.82 22.21 14 42.94 40.77 84.22 15 231.05 66.02 220.13 16 69.01 57.45 85.20 17 309.21 88.90 1394.35 18 13957.93 8632.00 35660.11
期望输出(前几行)
0 1 2 0 0.949 1.147 1.02 1 3.03 1.665 2.74 2 2.74 2.44 3.04 3 0.031 0.6 0.29
现有代码及问题
用户基于Dask编写的代码如下:
def create_dask(df): delayed_computations = [] # Convert each column to a Dask Series and then create a Dask DataFrame for column_name in df.columns: dask_series = dd.from_pandas(df[column_name], npartitions=1) #keep each column in one partition computation = dask_series.map_partitions(lambda df: calculate_ratios(df), meta=pd.Series([], dtype='float64')) delayed_computations.append(computation) return dask_results def calculate_ratios(data): # Calculate pairwise ratios for each value in the column file_len = 19 ratios = [] for i in range(file_len - 1): ratio_values = data.iloc[i] / data.iloc[i + 1:] ratios.append(ratio_values) return dd.concat(ratios, axis=1, ignore_index=True) if __name__ == '__main__': df_lung = pd.read_excel(os.path.join(base_path, 'raw_data_49_lung_cluster90 - Copy.xls')) dask_lung = create_dask(df_lung)
遇到的问题:原始Dask Series调用compute()正常,但后续computation调用compute()时执行极慢,无法验证函数逻辑。
问题原因分析
- Dask与Pandas操作混用:
calculate_ratios函数在map_partitions内部调用dd.concat,而map_partitions的输入是单个分区的Pandas Series,错误地在Pandas数据上执行Dask级别的拼接操作,生成大量不必要的Dask任务,导致性能暴跌。 - 硬编码数据长度:
file_len = 19是固定值,若数据行数变化会直接出错,且无法适配动态数据。 - 元数据不匹配:
map_partitions的meta参数设为pd.Series,但calculate_ratios实际返回的是多列DataFrame,元数据不匹配会导致Dask额外的类型推断或计算开销。 - 变量未定义:
create_dask函数最后返回的dask_results未定义,属于语法错误。 - 低效的列处理方式:逐个列转换为Dask Series再处理,不如直接用Dask DataFrame整体操作高效。
修复方案及代码
修复后的代码
import pandas as pd import dask.dataframe as dd import os def calculate_ratios(series): # 动态获取数据长度,替换硬编码 n = len(series) ratio_list = [] for i in range(n - 1): # 计算第i个值与后续所有值的比值,返回Series并重置索引 ratio = series.iloc[i] / series.iloc[i+1:] ratio_list.append(ratio.reset_index(drop=True)) # 将所有比值Series合并为DataFrame return pd.concat(ratio_list, axis=1, ignore_index=True) def create_dask(df): # 直接将整个DataFrame转为Dask DataFrame,保持单分区(列内全量计算需要完整数据) dask_df = dd.from_pandas(df, npartitions=1) # 用样本列计算结果生成meta模板,确保Dask能正确推断输出结构 sample_meta = calculate_ratios(df.iloc[:, 0]) # 按列应用计算函数并合并结果 dask_results = dask_df.map_partitions( lambda df: df.apply(calculate_ratios, axis=0), meta={col: sample_meta.dtypes for col in df.columns} ) return dask_results if __name__ == '__main__': base_path = "./" # 替换为你的实际路径 df_lung = pd.read_excel(os.path.join(base_path, 'raw_data_49_lung_cluster90 - Copy.xls')) dask_lung = create_dask(df_lung) # 验证前几行结果 print(dask_lung.head().round(3))
关键修复点
- 纯Pandas操作:
calculate_ratios仅处理Pandas Series,返回Pandas DataFrame,避免在分区内使用Dask操作,减少任务开销。 - 动态长度获取:用
len(series)替代硬编码的19,适配任意行数的数据。 - 正确的元数据设置:通过样本列计算结果生成meta模板,确保Dask能正确推断输出结构,避免额外计算。
- 整体DataFrame处理:直接将整个DataFrame转为Dask DataFrame,用
apply按列处理,代码更简洁高效。 - 修复变量未定义问题:正确生成并返回
dask_results。
性能优化补充
如果数据量极大,单分区处理仍有压力,可以考虑:
- 若列间计算相互独立,可将每列作为一个独立任务并行处理。
- 配置Dask的
distributed调度器,充分利用多核资源加速计算。
内容的提问来源于stack exchange,提问作者Mika Bell
相关产品推荐
相关产品推荐

