如何高效在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 loop | 2.16 s ± 82.3 ms |
| adjusting chunk size to 2d-array & map_blocks | 1.26 s ± 301 ms |
| adjusting chunk size to column size & map_blocks | 926 ms ± 31.9 |
| map over columns by slicing | 316 ms ± 11.5 ms |
| apply_along_axis | 1.01 s ± 18.7 ms |
| lazy loop | 1.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
相关产品推荐
相关产品推荐

