pandas数据处理遇VPN断连中断后如何实现断点续跑
Pandas长任务断点续跑实现方案
核心实现逻辑围绕进度自动持久化+启动自动续跑两个目标,不需要手动修改代码调整起始位置,同时可以灵活缩小分块大小降低断连带来的进度损失。
核心设计要点
- 用独立的轻量文件持久化记录最后一个处理完成的分块序号,仅当单个分块处理完成、结果成功写入本地磁盘后,才更新进度标记,避免进度记录和实际结果不一致
- 程序启动时自动扫描进度文件和已生成的分块结果,自动计算起始处理位置,无需手动改代码
- 取消提前切分全量DataFrame为分块列表的逻辑,按需按行号取分块,降低大表的内存占用
- 分块大小可根据单条数据处理耗时灵活调整(比如从500降到50/20),单次VPN断连最多损失1个分块的处理进度
- 最终合并结果时直接读取本地所有已存的分块csv,不依赖运行时内存里的列表,避免中途崩溃丢失内存数据
可直接复用的实现代码
import os import json import numpy as np import pandas as pd # ========== 配置参数,可按需调整 ========== CHUNK_SIZE = 50 # 单分块行数,调小可降低断连损失 PROGRESS_FILE = "process_progress.json" # 进度存储文件 RESULT_PREFIX = "processed_" # 分块结果文件名前缀 # ========================================== def load_progress(): """读取已存的进度,返回下一个要处理的分块序号、已完成的分块结果列表""" processed_chunks = [] start_chunk = 0 # 读取进度文件 if os.path.exists(PROGRESS_FILE): with open(PROGRESS_FILE, "r", encoding="utf-8") as f: progress = json.load(f) last_completed = progress.get("last_completed_chunk", -1) start_chunk = last_completed + 1 # 加载所有已完成的分块结果,同时校验文件完整性 for chunk_idx in range(start_chunk): chunk_path = f"{RESULT_PREFIX}{chunk_idx}.csv" if not os.path.exists(chunk_path): # 如果标记已完成但文件不存在,说明之前存写出错,从这个块重新跑 start_chunk = chunk_idx processed_chunks = processed_chunks[:chunk_idx] break try: chunk_df = pd.read_csv(chunk_path, index_col=0) processed_chunks.append(chunk_df) except Exception: # 文件损坏(比如写一半断连),从这个块重新跑 start_chunk = chunk_idx processed_chunks = processed_chunks[:chunk_idx] os.remove(chunk_path) break return start_chunk, processed_chunks def save_progress(chunk_idx): """保存当前处理完成的分块序号到进度文件""" with open(PROGRESS_FILE, "w", encoding="utf-8") as f: json.dump({"last_completed_chunk": chunk_idx}, f) if __name__ == "__main__": # 加载历史进度 current_chunk, together_list = load_progress() total_rows = len(large_df) total_chunks = int(np.ceil(total_rows / CHUNK_SIZE)) print(f"从分块{current_chunk}开始处理,总分块数{total_chunks}") # 从断点开始循环处理 while current_chunk < total_chunks: # 按行号取当前分块,不用提前切全量df start_row = current_chunk * CHUNK_SIZE end_row = min((current_chunk + 1) * CHUNK_SIZE, total_rows) chunk = large_df.iloc[start_row:end_row] try: # 处理分块(如果需要数据库连接重试,可在process_chunk内部加重试逻辑) chunk_processed = process_chunk(chunk) # 先存结果到本地 chunk_path = f"{RESULT_PREFIX}{current_chunk}.csv" chunk_processed.to_csv(chunk_path) # 存完结果再更新进度 save_progress(current_chunk) # 把结果加到内存列表 together_list.append(chunk_processed) print(f"分块{current_chunk}处理完成") current_chunk += 1 except Exception as e: # 遇到报错(比如VPN断连)直接终止循环,等用户重连后重启程序即可 print(f"分块{current_chunk}处理出错: {str(e)},程序退出,重连后重启将从当前断点继续") exit(1) # 全部处理完成后合并所有结果 all_chunks_together = pd.concat(together_list, ignore_index=True) all_chunks_together.to_csv("final_processed_result.csv", index=False) print("全量数据处理完成,最终结果已存为final_processed_result.csv")
使用说明
- VPN断连后只需要重新连上数据库,直接重新运行脚本即可,不需要修改任何代码,程序会自动找到上次中断的位置继续
- 如果需要调整分块大小,直接修改
CHUNK_SIZE参数即可,注意如果改小分块大小,建议先删掉旧的进度文件和已生成的processed_开头的csv,避免分块序号不匹配 - 可以在
process_chunk函数内部增加数据库连接失败的自动重试逻辑(比如捕获连接错误后等待30秒重连,重试3次再抛出异常),进一步减少人工重启的操作
内容的提问来源于stack exchange,提问作者Dr Wampa
相关产品推荐
相关产品推荐

