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

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()时执行极慢,无法验证函数逻辑。

问题原因分析

  1. Dask与Pandas操作混用:calculate_ratios函数在map_partitions内部调用dd.concat,而map_partitions的输入是单个分区的Pandas Series,错误地在Pandas数据上执行Dask级别的拼接操作,生成大量不必要的Dask任务,导致性能暴跌。
  2. 硬编码数据长度:file_len = 19是固定值,若数据行数变化会直接出错,且无法适配动态数据。
  3. 元数据不匹配:map_partitions的meta参数设为pd.Series,但calculate_ratios实际返回的是多列DataFrame,元数据不匹配会导致Dask额外的类型推断或计算开销。
  4. 变量未定义:create_dask函数最后返回的dask_results未定义,属于语法错误。
  5. 低效的列处理方式:逐个列转换为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))

关键修复点

  1. 纯Pandas操作:calculate_ratios仅处理Pandas Series,返回Pandas DataFrame,避免在分区内使用Dask操作,减少任务开销。
  2. 动态长度获取:用len(series)替代硬编码的19,适配任意行数的数据。
  3. 正确的元数据设置:通过样本列计算结果生成meta模板,确保Dask能正确推断输出结构,避免额外计算。
  4. 整体DataFrame处理:直接将整个DataFrame转为Dask DataFrame,用apply按列处理,代码更简洁高效。
  5. 修复变量未定义问题:正确生成并返回dask_results。

性能优化补充

如果数据量极大,单分区处理仍有压力,可以考虑:

  • 若列间计算相互独立,可将每列作为一个独立任务并行处理。
  • 配置Dask的distributed调度器,充分利用多核资源加速计算。

内容的提问来源于stack exchange,提问作者Mika Bell

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 06:42:06