Dask易并行for循环优化:避免DataFrame复制以提速计算
优化方案:跳过DataFrame复制,降低迭代耗时
你的核心瓶颈是每次迭代复制150万行的DataFrame,其实完全可以避免——因为每个迭代只用到原df的important和ID列,且最终只需要返回聚合后的结果,不需要修改原数据。以下是具体优化思路和代码:
1. 提前提取核心数据,避免全量复制
先从原DataFrame中提取计算必需的列转为numpy数组(比DataFrame复制成本低几个数量级):
import numpy as np from scipy.stats import norm import pandas as pd # 提前提取核心数据,仅执行一次 important_arr = df["important"].values id_arr = df["ID"].values
2. 修改迭代函数,完全跳过DataFrame复制
直接用提取的数组完成计算,最后仅构造小量级的聚合结果返回,全程不复制原DataFrame:
def iteration(important_arr, id_arr, seed): rng = np.random.default_rng(seed) factor = norm.ppf(rng.uniform()) factors = norm.ppf(rng.uniform(size=len(important_arr))) # 直接用数组计算,无DataFrame复制操作 result_arr = important_arr * factor + important_arr * factors # 构造小结果DataFrame并完成聚合 return pd.DataFrame({ "ID": id_arr, "result": result_arr }).groupby(by="ID").agg({"result": "sum"}).reset_index()
3. 保持并行逻辑不变,传入轻量数据
原Dask并行逻辑无需大改,仅将传入参数替换为提前提取的轻量数组:
sg = np.random.SeedSequence(112622607388004198928332989692281456100) seeds = sg.generate_state(10000) delayed_results = [] for seed in seeds: # 传入轻量数组而非整个DataFrame result = dask.delayed(iteration)(important_arr, id_arr, seed) delayed_results.append(result) results = dask.compute(*delayed_results)
额外提速:用numpy替代pandas分组
如果想进一步压缩聚合耗时,可以用numpy原生方法实现分组求和,速度比pandas的groupby更快:
def iteration(important_arr, id_arr, seed): rng = np.random.default_rng(seed) factor = norm.ppf(rng.uniform()) factors = norm.ppf(rng.uniform(size=len(important_arr))) result_arr = important_arr * factor + important_arr * factors # numpy原生分组求和,性能更优 unique_ids, idx = np.unique(id_arr, return_inverse=True) group_sums = np.bincount(idx, weights=result_arr) return pd.DataFrame({"ID": unique_ids, "result": group_sums})
这样修改后,每次迭代不再复制百万行的DataFrame,仅处理轻量的numpy数组,耗时会大幅降低;同时通过SeedSequence生成的种子保证了结果可复现,也彻底避免了之前多列插入报错的问题。
内容的提问来源于stack exchange,提问作者bogdmu00
相关产品推荐
相关产品推荐

