如何让逐行循环适配Pandas/Modin/Ray并避免内存溢出
优化Pandas/Modin循环避免内存溢出并生成目标字典
问题背景
我有一个需逐行执行的半复杂循环,尝试过相关资料但仍无法理解如何借助Pandas/Modin/Ray生成目标字典。当前循环在12万行、每行带大型列表的数据集上运行时,即便有100GB内存仍会内存溢出导致Python崩溃,需要修改代码避免该问题。
当前简化代码
import modin.pandas as pd import numpy as np import ray ray.init() df = pd.read_csv("file.csv") dict_res = {} for index, row in df.iterrows(): list_items = row['listed_items'] length = len(list_items) for i in range(0, length): for j in range(i+1, length): key_str = "{},{}".format(list_items[i], list_items[j]) if key_str in dict_res: dict_res[key_str] += (1/(length-1)) else: dict_res[key_str] = (1/(length-1))
数据与结果示例
输入数据示例
row1 = [100000, 200000, 421563] row2 = [500, 453100, 442211, ...]
目标结果示例
dict_res = { "100000,200000" : 0.5, "100000,421563" : 0.5, "200000,421563" : 0.5, ... }
测试用例
testfile.csv内容
prop,items XY108,"[9929, 102010, 301352, 521008]" XY109,"[382, 396, 456, 639, 883, 1291, 1333, 1969, 9929, 102010, 11457, 12425, 15770]"
生成测试DataFrame的代码
from collections import Counter import modin.pandas as pd import ray ray.init() df = pd.read_csv("testfile.csv") def str_to_list(list_str): return [int(x) for x in list_str.strip('[]').split(',')] df['items'] = df['items'].apply(str_to_list)
优化思路与解决方案
核心问题分析
iterrows()逐行处理效率极低,且会将全量数据加载到内存;- 嵌套循环生成大量重复/冗余键值对,字典内存占用随数据量呈指数增长;
- 全局字典的频繁修改会引发锁竞争,分布式环境下效率进一步下降。
优化方案
1. 向量化分组聚合(推荐)
利用Pandas/Modin的向量化操作,先生成所有合法元素对,再按分组聚合计算权重,避免嵌套循环:
import modin.pandas as pd from itertools import combinations import ray ray.init() # 读取并预处理数据 df = pd.read_csv("testfile.csv") df['items'] = df['items'].apply(lambda x: [int(i) for i in x.strip('[]').split(',')]) # 生成单条数据的元素对及对应权重 def generate_pairs(row): items = row['items'] n = len(items) if n < 2: return [] weight = 1/(n-1) # 排序元素避免生成"a,b"和"b,a"重复键 return [(f"{a},{b}", weight) for a, b in combinations(sorted(items), 2)] # 展开所有元素对并聚合求和 pair_series = df.apply(generate_pairs, axis=1).explode().dropna() pair_df = pair_series.apply(pd.Series, index=['key', 'weight']) dict_res = pair_df.groupby('key')['weight'].sum().to_dict() print(dict_res)
2. 分布式分块处理(超大数据集适用)
通过Ray将数据分块,分布式计算每个块的结果后再合并,避免单进程内存过载:
import modin.pandas as pd from itertools import combinations import ray from collections import defaultdict import numpy as np ray.init() @ray.remote def process_chunk(chunk): chunk['items'] = chunk['items'].apply(lambda x: [int(i) for i in x.strip('[]').split(',')]) chunk_res = defaultdict(float) for _, row in chunk.iterrows(): items = row['items'] n = len(items) if n < 2: continue weight = 1/(n-1) for a, b in combinations(sorted(items), 2): chunk_res[f"{a},{b}"] += weight return chunk_res # 读取数据并分块(块数量可根据内存调整) df = pd.read_csv("testfile.csv") chunks = np.array_split(df, 10) # 提交分布式任务并合并结果 futures = [process_chunk.remote(chunk) for chunk in chunks] results = ray.get(futures) dict_res = defaultdict(float) for res in results: for key, val in res.items(): dict_res[key] += val dict_res = dict(dict_res) print(dict_res)
3. 关键内存优化细节
- 对元素对排序,避免生成"a,b"和"b,a"两种重复键,直接减少字典键的数量;
- 使用
defaultdict(float)替代普通字典,省去键存在性检查的开销; - 分块处理时,每个块单独计算后再合并,避免一次性加载所有元素对到内存。
内容的提问来源于stack exchange,提问作者Ranger
相关产品推荐
相关产品推荐

