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

如何无循环处理大型Pandas DataFrame行存入列表(解决内存溢出)

内存友好的大型DataFrame行段筛选方案

针对你500万行DataFrame的需求,以下是几个避免内存耗尽的具体实现方法,核心思路是用矢量化操作替代嵌套循环、分批处理/存储,避免一次性加载所有结果到内存:

方案1:Pandas矢量化标记+分批存储

先通过矢量化操作标记每个行程内的有效连续段,再分批将结果写入文件(或拆分到多个小列表),完全规避嵌套循环的低效和内存堆积:

import pandas as pd

# 1. 标记有效行并生成连续段ID
df['valid'] = df['cum'] >= 10
# 给每个trip内的连续有效段分配唯一ID
df['segment_id'] = df.groupby('trip')['valid'].apply(
    lambda x: (x != x.shift()).cumsum()
)
# 过滤出仅包含有效段的行
valid_segments = df[df['valid']].copy()

# 2. 分批处理并存储
batch_size = 1000  # 可根据内存调整批次大小
current_batch = []
batch_num = 1

# 按(trip, segment_id)遍历每个连续行段
for (trip_id, seg_id), segment in valid_segments.groupby(['trip', 'segment_id']):
    # 转换为你需要的格式(字典列表/值列表,按需选择)
    seg_data = segment.values.tolist()  # 比to_dict更省内存
    current_batch.append(seg_data)
    
    # 达到批次阈值时保存并清空当前批次
    if len(current_batch) >= batch_size:
        pd.to_pickle(current_batch, f'./segment_batch_{batch_num}.pkl')
        # 若需内存中保留多个小列表,可改为 all_batches.append(current_batch)
        current_batch = []
        batch_num += 1

# 处理剩余的未达批次的行段
if current_batch:
    pd.to_pickle(current_batch, f'./segment_batch_{batch_num}.pkl')

方案2:用Dask分块处理超大数据集

如果DataFrame大到单Pandas无法处理,用Dask将数据拆分为多个分区并行处理,全程仅加载部分数据到内存:

import dask.dataframe as dd

# 将Pandas DataFrame转为Dask DataFrame,分区数根据内存调整
ddf = dd.from_pandas(df, npartitions=8)

# 定义每个分区的处理函数
def process_partition(partition):
    partition['valid'] = partition['cum'] >= 10
    partition['segment_id'] = partition.groupby('trip')['valid'].apply(
        lambda x: (x != x.shift()).cumsum()
    )
    valid_segs = partition[partition['valid']]
    # 返回当前分区内的所有有效行段列表
    return valid_segs.groupby(['trip', 'segment_id']).apply(lambda x: x.values.tolist()).tolist()

# 分块计算结果(仅加载部分分区到内存)
partition_results = ddf.map_partitions(process_partition).compute()

# 合并分区结果并分批存储
batch_size = 1000
current_batch = []
batch_num = 1
for seg_list in partition_results:
    current_batch.extend(seg_list)
    if len(current_batch) >= batch_size:
        pd.to_pickle(current_batch, f'./dask_segment_batch_{batch_num}.pkl')
        current_batch = []
        batch_num += 1
if current_batch:
    pd.to_pickle(current_batch, f'./dask_segment_batch_{batch_num}.pkl')

方案3:生成器按需迭代(内存占用最低)

用生成器函数逐个生成有效行段,全程仅在内存中保留单个行段的数据,彻底避免内存堆积:

def generate_valid_segments(df):
    df['valid'] = df['cum'] >= 10
    df['segment_id'] = df.groupby('trip')['valid'].apply(
        lambda x: (x != x.shift()).cumsum()
    )
    valid_segments = df[df['valid']]
    # 逐个生成有效行段
    for (trip_id, seg_id), segment in valid_segments.groupby(['trip', 'segment_id']):
        yield segment.values.tolist()

# 遍历生成器并分批存储
batch_size = 1000
current_batch = []
batch_num = 1
for seg_data in generate_valid_segments(df):
    current_batch.append(seg_data)
    if len(current_batch) >= batch_size:
        pd.to_pickle(current_batch, f'./generator_batch_{batch_num}.pkl')
        current_batch = []
        batch_num += 1
if current_batch:
    pd.to_pickle(current_batch, f'./generator_batch_{batch_num}.pkl')

关键优化点

  • 避免嵌套循环:用Pandas的groupby+矢量化操作替代手动循环,效率提升数倍且内存占用更低。
  • 选择紧凑存储格式:用values.tolist()代替to_dict('records'),减少内存开销;用Pickle保存批次比CSV更高效。
  • 优先写入文件:若无需全程内存保留结果,直接分批写入文件是最省内存的方式。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 14:35:56