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

如何让逐行循环适配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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 05:08:10