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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 22:03:23