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

Dask字符串DataFrame逐行apply过慢,求优化方案

问题:优化Dask DataFrame批量字符串替换效率

需求:

  • 处理无缺失值的Dask DataFrame,对除前两列(R、A)外的所有列执行以下操作:
    1. 目标列均为二进制字符串,将其中的0替换为对应行的R值,1替换为对应行的A值;
    2. 将替换后的字符串转为小写;
    3. 最终删除R、A列。

示例输入:

R  A   T   U   V
0   R  A  00  10  11
.. .. ..  ..  ..  ..
95  R  A  11  00  00

示例输出:

T   U   V
0   rr  ar  aa
..  ..  ..  ..
95  aa  rr  rr

原代码(运行速度过慢):

from dask.distributed import Client, progress
client = Client(n_workers=20, threads_per_worker=1)
client

import pandas as pd
import dask.dataframe as dd

def score(x):
    return str(x[0] + x[1]).lower()

def func(row, i, fs):
    r = row[0]
    a = row[1]
    row[i] = fs(row[i].replace('0',r).replace('1',a))
    return row

s1 = ['00']*25 + ['01']*25 + ['10']*25 + ['11']*25
s2 = ['10']*25 + ['11']*25 + ['01']*25 + ['00']*25

df = pd.DataFrame({'R':['R']*100, 'A':['A']*100, 'T':s1, 'U':s2, 'V': reversed(s1)})
ddf = dd.from_pandas(df, npartitions=10)

meta = dict()
for cn in ddf.columns:
    meta[cn] = 'object'

ddf.compute()
for i in range(2,len(ddf.columns)):
    ddf = ddf.apply(func, args=(i,score,), axis=1, meta=meta)

ddf = ddf.drop(['R','A'], axis=1)
ddf.compute()

优化方案

原代码性能瓶颈分析

  1. 逐行apply操作:Dask中apply(axis=1)是逐行处理,无法利用向量化优化,调度和执行开销极大,尤其当数据量较大时。
  2. 循环处理列:多次调用apply生成多个计算步骤,增加了计算图复杂度和调度次数。
  3. 不必要的compute调用:提前执行ddf.compute()会打断Dask的延迟执行机制,失去分区并行优化的优势。

优化代码实现

利用Dask的向量化字符串操作,一次性批量处理所有目标列:

from dask.distributed import Client, progress
import pandas as pd
import dask.dataframe as dd

# 初始化客户端(按需调整参数)
client = Client(n_workers=20, threads_per_worker=1)

# 生成测试数据
s1 = ['00']*25 + ['01']*25 + ['10']*25 + ['11']*25
s2 = ['10']*25 + ['11']*25 + ['01']*25 + ['00']*25
df = pd.DataFrame({'R':['R']*100, 'A':['A']*100, 'T':s1, 'U':s2, 'V': list(reversed(s1))})
ddf = dd.from_pandas(df, npartitions=10)

# 获取目标列(除R、A外的所有列)
target_cols = ddf.columns[2:]

# 定义向量化处理函数:对单个列执行替换+转小写
def process_col(col_series, r_series, a_series):
    # 先替换0为R值,再替换1为A值,最后转小写
    return col_series.str.replace('0', r_series, regex=False)\
                     .str.replace('1', a_series, regex=False)\
                     .str.lower()

# 批量处理所有目标列
processed_cols = {col: process_col(ddf[col], ddf['R'], ddf['A']) for col in target_cols}

# 构造新的DataFrame,只保留处理后的列
ddf_processed = dd.concat([processed_cols[col] for col in target_cols], axis=1)

# 执行计算(仅在需要结果时调用)
result = ddf_processed.compute()
print(result)

优化点说明

  1. 向量化操作:使用Dask的str.replace和str.lower向量化方法,替代逐行apply,充分利用Dask的并行计算能力和Pandas的向量化优化。
  2. 批量处理列:一次性生成所有处理后的列,避免循环调用apply,简化计算图。
  3. 保留延迟执行:仅在最后需要结果时调用compute,让Dask自动优化计算流程和分区并行。
  4. 避免修改原DataFrame:通过构造新的列集合生成结果,减少不必要的数据拷贝和修改操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 22:43:18