如何加速100M行×10K列矩阵的Dask行过滤?
优化Dask大行列数据的行级过滤性能
原代码的核心问题是逐行使用Python的Counter进行统计,对于10K列的行来说,Counter的纯Python遍历开销极大,且Dask的行级apply无法充分利用列分区的并行优化优势。以下是针对性的改进方案:
1. 替换Counter为Numpy的向量化统计
Numpy的底层实现是C级别的,比Python原生的Counter快数倍。修改retention函数,用numpy的unique替代Counter:
import numpy as np def retention(row): min_retention = 90.0 arr = row.to_numpy() total_cols = len(arr) if total_cols == 0: return False # 统计唯一值及出现次数 _, counts = np.unique(arr, return_counts=True) max_count = counts.max() return (100.0 * max_count / total_cols) > min_retention fil_df = df[df.apply(retention, axis=1, meta=(None, 'bool'))]
2. 改用Dask Array进行底层并行优化
Dask DataFrame的行级apply调度开销较大,转为Dask Array后,利用map_blocks结合np.apply_along_axis可以获得更好的并行性能:
import dask.array as da # 将DataFrame转为Dask Array(lengths=True确保行长度一致) darr = df.to_dask_array(lengths=True) total_cols = darr.shape[1] def calc_max_count(arr): # 对每行计算最高频值的出现次数 return np.apply_along_axis(lambda x: np.unique(x, return_counts=True)[1].max(), 1, arr) # 计算每行的最高频计数 max_counts = da.map_blocks(calc_max_count, darr, dtype=np.int64) # 生成过滤掩码 retention_mask = (100.0 * max_counts / total_cols) > 90.0 # 应用掩码到原DataFrame fil_df = df[retention_mask]
3. 调整Dask分区与调度器
- 优化分区大小:默认分区可能过小导致调度开销过大,根据内存和CPU核心数调整分区,比如让每个分区占1-2GB内存:
df = df.repartition(npartitions=64) # 示例值,根据实际资源调整 - 使用进程/分布式调度器:对于CPU密集型任务,进程调度器比线程调度器更高效,或者使用Dask分布式集群进一步提升并行能力:
import dask dask.config.set(scheduler='processes') # 多进程调度
4. 特殊场景下的近似过滤(可选)
如果业务允许一定精度误差,可以通过随机采样列减少计算量,比如每行只采样10%的列来估算占比:
def retention_sampled(row): min_retention = 90.0 # 随机采样1000列(10K的10%) sampled = row.sample(n=1000, random_state=42).to_numpy() _, counts = np.unique(sampled, return_counts=True) max_count = counts.max() return (100.0 * max_count / len(sampled)) > min_retention fil_df = df[df.apply(retention_sampled, axis=1, meta=(None, 'bool'))]
内容的提问来源于stack exchange,提问作者xinit
相关产品推荐
相关产品推荐

