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

Pandas大数据聚合内存不足:分块处理与内存优化方案咨询

问题解决与大数据集优化方案

直接错误修复

你遇到的MemoryError是代码笔误导致的:在交易特征循环中,错误使用transaction_fields(所有交易列的列表)代替transaction_field(当前循环的单个列),触发了广播机制——用列列表除以单个列时,Pandas会将单个列扩展为与列列表行数一致的二维数组,最终生成(208719, 208719)的超大矩阵,直接超出内存上限。

修正后的交易特征计算代码:

# Sum transactions
transaction_fields = [column for column in data.columns if 'Transactions' in column and 'SavingAccount' in column]
# 用向量sum代替逐行apply,更高效
data['Transactions'] = data[transaction_fields].sum(axis=1)
for transaction_field in transaction_fields:
    data[f'Transactions_{transaction_field}_vs_total'] = np.where(
        data['Transactions'] == 0,
        0,  # 直接用0代替floor(0),结果一致
        np.floor(data[transaction_field] / data['Transactions'] * 100)
    )

同时建议优化之前的Operations计算,改用向量运算:

# Sum Operations 优化版
operation_fields = [column for column in data.columns if 'Operations' in column]
data['Operations'] = data[operation_fields].sum(axis=1)
for operation_field in operation_fields:
    data[f'Operations_{operation_field}_vs_total'] = np.where(
        data['Operations'] == 0,
        0,
        np.floor(data[operation_field] / data['Operations'] * 100)
    )

内存优化策略

1. 优先使用向量运算

apply(axis=1)是逐行遍历,效率低且内存开销大;Pandas内置的向量运算(如sum(axis=1)、除法、乘法)基于Numpy实现,速度快且内存占用少,尽量避免逐行处理。

2. 压缩数据类型

  • 检查数值列的取值范围,将float64转成float32(精度允许的情况下),或用整数类型存储百分比结果(比如把floor后的结果转成int8或int16):
# 转换Transactions列类型
data['Transactions'] = data['Transactions'].astype('int32')
# 转换百分比列类型
for col in data.columns:
    if 'vs_total' in col:
        data[col] = data[col].astype('int8')

3. 清理冗余数据

  • 计算完聚合特征后,删除不再需要的原始列,释放内存:
data.drop(columns=operation_fields + transaction_fields, inplace=True)
  • 删除不再使用的变量,手动触发垃圾回收:
import gc
del operation_fields, transaction_fields
gc.collect()

4. 分块处理方案

如果数据集过大无法一次性加载,用Pandas分块读取功能:

# 分块读取原始数据(假设数据来自CSV)
chunk_size = 10000
chunks = []
for chunk in pd.read_csv('your_data.csv', chunksize=chunk_size):
    # 在每个块上执行特征计算
    operation_fields = [col for col in chunk.columns if 'Operations' in col]
    chunk['Operations'] = chunk[operation_fields].sum(axis=1)
    for op_field in operation_fields:
        chunk[f'Operations_{op_field}_vs_total'] = np.where(
            chunk['Operations'] == 0,
            0,
            np.floor(chunk[op_field] / chunk['Operations'] * 100)
        )
    
    transaction_fields = [col for col in chunk.columns if 'Transactions' in col and 'SavingAccount' in col]
    chunk['Transactions'] = chunk[transaction_fields].sum(axis=1)
    for trans_field in transaction_fields:
        chunk[f'Transactions_{trans_field}_vs_total'] = np.where(
            chunk['Transactions'] == 0,
            0,
            np.floor(chunk[trans_field] / chunk['Transactions'] * 100)
        )
    
    # 删除冗余列
    chunk.drop(columns=operation_fields + transaction_fields, inplace=True)
    chunks.append(chunk)

# 合并所有块
final_data = pd.concat(chunks, ignore_index=True)

5. 用Dask替代Pandas

Dask是专门处理超大数据集的工具,语法与Pandas几乎兼容,会自动分块并行计算,无需手动处理分块:

import dask.dataframe as dd

# 用Dask加载数据
dask_data = dd.read_csv('your_data.csv')

# 执行特征计算(语法和Pandas一致)
operation_fields = [col for col in dask_data.columns if 'Operations' in col]
dask_data['Operations'] = dask_data[operation_fields].sum(axis=1)
for op_field in operation_fields:
    dask_data[f'Operations_{op_field}_vs_total'] = np.where(
        dask_data['Operations'] == 0,
        0,
        np.floor(dask_data[op_field] / dask_data['Operations'] * 100)
    )

# 计算并获取结果(或直接保存到文件)
final_data = dask_data.compute()

内容的提问来源于stack exchange,提问作者FedeG

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 16:43:13