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

如何高效在Dask数组列上应用函数?基准测试解析

问题描述

我有一个尺寸约为3,000,000×10,000的大型Dask数组,需要对每一列应用ECDF函数,最终结果维度需与原数组一致。我测试了多种实现方法并完成基准测试,发现按列切片映射速度最快,但这与我对Dask的预期不符,怀疑操作存在问题,希望了解最优方案及改进方向。

测试代码

import numpy as np
import dask
import dask.array as da
from dask.distributed import Client, LocalCluster
from statsmodels.distributions.empirical_distribution import ECDF

### functions
def ecdf(x):
    fn = ECDF(x)
    return fn(x)

def ecdf_array(X):
    res = np.zeros_like(X)
    for i in range(X.shape[1]):
        res[:,i] = ecdf(X[:,i])
    return res

### set up scheduler / workers
n_workers = 10
cluster = LocalCluster(n_workers=n_workers, threads_per_worker=4)
client = Client(cluster)

### create data set 
X = da.random.random((100000,100)) #dask
Xarr = X.compute() #numpy

### traditional for loop
%timeit -r 10 foo1 = ecdf_array(Xarr)

### adjusting chunk size to 2d-array and map_blocks
X = X.rechunk(chunks=(X.shape[0],np.ceil(X.shape[1]/n_workers)))
Xm = X.map_blocks(lambda x: ecdf_array(x),meta = np.array((), dtype='float'))
%timeit -r 10 foo2 = Xm.compute()

### adjusting chunk size to column size and map_blocks
X = X.rechunk(chunks=(X.shape[0],1))
Xm = X.map_blocks(lambda x: np.expand_dims(ecdf(np.squeeze(x)),1),meta = np.array((), dtype='float'))
%timeit -r 10 foo3 = Xm.compute()

### map over columns by slicing
Xm = client.map(lambda i: ecdf(np.asarray(X[:,i])),range(X.shape[1]))
Xm = client.submit(lambda x: da.transpose(da.vstack(x)),Xm)
%timeit -r 10 foo4 = Xm.result()

### apply_along_axis
Xaa = da.apply_along_axis(lambda x: np.expand_dims(ecdf(x),1), 0, X, dtype=X.dtype, shape=X.shape)
%timeit -r 10 foo5 = Xaa.compute()

### lazy loop
Xl = []
for i in range(X.shape[1]):
    Xl.append(dask.delayed(ecdf)(X[:,i]))
Xl = dask.delayed(da.vstack)(Xl)
%timeit -r 10 foo6 = Xl.compute()

基准测试结果

方法Results (10 loops)
traditional for loop2.16 s ± 82.3 ms
adjusting chunk size to 2d-array & map_blocks1.26 s ± 301 ms
adjusting chunk size to column size & map_blocks926 ms ± 31.9
map over columns by slicing316 ms ± 11.5 ms
apply_along_axis1.01 s ± 18.7 ms
lazy loop1.4 s ± 352 ms

解答

为什么预期的方法效率低

你预期的“调整块大小为二维数组+map_blocks”方法慢的核心原因是没有利用worker的多线程资源:

  • 该方法将多列打包到一个chunk中,每个worker处理一个chunk,但ecdf_array内部是单线程循环遍历列,相当于每个worker只用到了1个线程(你设置了threads_per_worker=4),浪费了3/4的并行能力。
  • ECDF计算是列独立的操作,多列打包的chunk无法让worker的多个线程同时处理不同列,反而因为单线程循环拖慢了整体速度。

为什么“按列切片映射”最快

该方法的优势在于完全利用了集群的并行资源:

  • client.map直接将每一列的ECDF计算任务分发到集群的所有线程(10个worker×4线程=40个并行单元),每个任务都是独立的,没有单线程循环的瓶颈。
  • 任务粒度(单列)与ECDF的计算粒度完全匹配,调度开销相对可控,最大化了并行效率。

最优方案与改进方向

1. 优化ECDF计算本身

statsmodels的ECDF封装了额外的类逻辑,用numpy原生操作实现可以大幅提速:

def fast_ecdf(x):
    x_sorted = np.sort(x)
    # 用searchsorted计算每个元素在排序数组中的位置,得到ECDF值
    return np.searchsorted(x_sorted, x, side='right') / len(x_sorted)

这个实现避免了类实例化的开销,比statsmodels版本快2-3倍。

2. 平衡任务粒度与调度开销

单列并行虽然快,但10,000列会产生10,000个小任务,超大规模下调度开销会上升。可以将若干列打包成批次,平衡并行度和调度成本:

def ecdf_batch(col_indices):
    batch_results = []
    for i in col_indices:
        col = X[:, i].compute()
        batch_results.append(fast_ecdf(col))
    return da.vstack(batch_results)

# 按每10列分批次,可根据集群规模调整
batch_size = 10
batches = [range(i, min(i+batch_size, X.shape[1])) for i in range(0, X.shape[1], batch_size)]
batch_tasks = client.map(ecdf_batch, batches)
final_result = da.transpose(da.hstack(client.gather(batch_tasks)))

3. 优化Dask数组的chunk结构

确保原数组的chunk是行完整、列拆分的结构(即chunks=(3_000_000, 1)或合适的列批次大小),这样X[:,i]的切片操作无需跨chunk合并,减少数据传输开销。

4. 验证结果一致性

在优化前,务必对比fast_ecdf与statsmodels ECDF的输出结果,确保计算逻辑一致(可随机选取若干列进行数值对比)。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 00:01:22