如何无循环处理大型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
相关产品推荐
相关产品推荐

