如何在Pandas DataFrame的apply处理中每N步自动保存数据?
大DataFrame分批处理并自动保存的优雅实现
问题场景
我的DataFrame数据量极大,想找一种不用for循环的优雅方法,在修改DataFrame内部分值的同时,每N步执行一次保存操作。现有示例代码如下:
def modifier(x): x = x.split() # 此处会应用更复杂的逻辑 return x df['new_col'] = df.old_col.apply(modifier)
能否在modifier函数中添加代码,使其每处理10000行数据时,自动调用df.to_pickle('make_copy.pickle')执行保存?
可行方案
直接在modifier里嵌入保存逻辑并不靠谱——apply是逐行处理,函数内部没法准确追踪处理进度,而且刚处理的行还没同步到DataFrame里,此时保存会导致数据不完整。推荐两种更稳妥的实现方式:
1. 计数器+apply(轻量单机方案)
维护一个全局计数器,每处理指定行数就触发保存,配合tqdm还能看到处理进度:
from tqdm import tqdm import pandas as pd # 配置参数 counter = 0 BATCH_SIZE = 10000 def modifier(x): global counter counter += 1 # 替换成你的复杂处理逻辑 result = x.split() # 每处理BATCH_SIZE行自动保存 if counter % BATCH_SIZE == 0: df.to_pickle('make_copy.pickle') return result # 带进度条的apply执行 tqdm.pandas(desc="处理数据中") df['new_col'] = df.old_col.progress_apply(modifier) # 最后处理剩余不足BATCH_SIZE的行,补一次保存 df.to_pickle('make_copy.pickle')
注意:全局计数器仅适用于单进程处理,多线程/多环境下可能出现计数错误。
2. 分块处理(超大数据集首选)
如果DataFrame大到内存无法承载,直接分块读取处理,每块处理完成后保存,最后合并结果:
import pandas as pd # 配置参数 BATCH_SIZE = 10000 final_output = 'final_result.pickle' temp_prefix = 'temp_batch' # 分块读取原数据(如果是csv可以用read_csv,pickle用read_pickle) for batch_idx, chunk in enumerate(pd.read_pickle('original_data.pickle', chunksize=BATCH_SIZE)): # 处理当前块 chunk['new_col'] = chunk.old_col.apply(lambda x: x.split()) # 替换为你的modifier逻辑 # 保存临时块 chunk.to_pickle(f'{temp_prefix}_{batch_idx}.pickle') # 每处理5个小批次合并一次总结果(可根据需求调整) if (batch_idx + 1) % 5 == 0: batch_list = [pd.read_pickle(f'{temp_prefix}_{i}.pickle') for i in range(batch_idx-3, batch_idx+1)] combined_df = pd.concat(batch_list) combined_df.to_pickle(final_output) # 可选:删除临时文件释放空间 # 最后合并所有剩余批次 all_batches = [pd.read_pickle(f'{temp_prefix}_{i}.pickle') for i in range(batch_idx+1)] final_df = pd.concat(all_batches) final_df.to_pickle(final_output)
这种方式避免了一次性加载全量数据,而且每一步都有保存,就算程序中途崩溃也不会丢失全部进度。
为什么不推荐在modifier里直接保存?
apply的执行顺序不一定严格对应行号,计数器可能出现偏差- 函数内无法确保当前处理的行已经写入DataFrame,保存的内容可能不完整
- 逐行判断是否保存会额外增加性能开销,不如按固定间隔或分块保存高效
内容的提问来源于stack exchange,提问作者TommyLeeJones
相关产品推荐
相关产品推荐

