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

Windows10下无管理员权限的批量CSV中断续处理方案咨询

高效实现断点续跑与单条追加的CSV批量处理方案

核心设计思路

  • 用本地JSON文件记录处理进度,包括已完成的文件、每个文件已处理的行数
  • 单条/分块记录追加输出到结果CSV,避免批量写入导致的进度丢失
  • 基于进度文件实现断点恢复,重启后自动跳过已处理内容
  • 针对高负载NLP处理做内存与效率优化

1. 进度追踪模块

创建progress.json存储处理状态,脚本启动时先读取该文件确定起始位置:

import json
import os
import glob
import pandas as pd
import logging

# 基础配置
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
PROGRESS_FILE = "progress.json"
OUTPUT_CSV = "processed_results.csv"
FILES_PATH = "your_csv_folder/*.csv"  # 替换为实际文件夹路径

def load_progress():
    """加载已处理进度"""
    if os.path.exists(PROGRESS_FILE):
        with open(PROGRESS_FILE, 'r', encoding='utf-8') as f:
            return json.load(f)
    # 初始进度模板
    return {"completed_files": [], "current_file": None, "current_row": 0}

def save_progress(progress):
    """实时保存当前进度"""
    with open(PROGRESS_FILE, 'w', encoding='utf-8') as f:
        json.dump(progress, f, indent=2)

2. 断点恢复与单条追加处理

修改原有逻辑,加入进度检查,分块处理CSV并实时追加结果:

def clean_process(row):
    """你的自定义NLP清洗逻辑(改为处理单条记录)"""
    # 替换为你的实际代码:比如正则清洗、实体识别等
    # row['cleaned_text'] = 你的处理逻辑
    return row

def process_single_file(file_path, start_row=0):
    """处理单个CSV文件,从指定行开始,逐块追加结果"""
    # 判断是否需要写入表头
    write_header = not os.path.exists(OUTPUT_CSV)
    # 分块读取(chunk_size可根据内存调整)
    chunk_size = 100
    
    for chunk_idx, chunk in enumerate(pd.read_csv(file_path, chunksize=chunk_size)):
        chunk_start = chunk_idx * chunk_size
        # 跳过已处理的块
        if chunk_start + len(chunk) <= start_row:
            continue
        # 截取当前需要处理的行
        if chunk_start < start_row:
            chunk = chunk.iloc[start_row - chunk_start:]
        
        # 处理当前块的每一行
        processed_chunk = chunk.apply(clean_process, axis=1)
        # 过滤处理失败的行(可选)
        processed_chunk = processed_chunk.dropna()
        
        # 追加到输出CSV
        processed_chunk.to_csv(OUTPUT_CSV, mode='a', header=write_header, index=False, encoding='utf-8')
        write_header = False
        
        # 实时更新进度
        current_row = chunk_start + len(chunk)
        progress = load_progress()
        progress["current_file"] = file_path
        progress["current_row"] = current_row
        save_progress(progress)
        
        logging.info(f"{file_path} 已处理至第 {current_row} 行")
    
    # 标记当前文件处理完成
    progress = load_progress()
    if file_path not in progress["completed_files"]:
        progress["completed_files"].append(file_path)
    progress["current_file"] = None
    progress["current_row"] = 0
    save_progress(progress)
    logging.info(f"{file_path} 处理完成")

def main():
    progress = load_progress()
    completed_files = progress["completed_files"]
    current_file = progress["current_file"]
    current_row = progress["current_row"]
    
    # 获取所有待处理文件(排序保证处理顺序一致)
    all_files = sorted(glob.glob(FILES_PATH))
    
    # 恢复处理未完成的文件
    if current_file and current_file in all_files:
        logging.info(f"恢复处理:{current_file},从第 {current_row} 行开始")
        process_single_file(current_file, current_row)
    
    # 处理剩余未完成的文件
    for file in all_files:
        if file not in completed_files:
            logging.info(f"开始处理:{file}")
            process_single_file(file)
    
    # 所有文件处理完成后删除进度文件(可选)
    if os.path.exists(PROGRESS_FILE):
        os.remove(PROGRESS_FILE)
    logging.info("全部文件处理完成!")

3. 自动重启实现(无管理员权限)

用批处理脚本实现自动重启,无需系统权限:

创建run_process.bat文件,内容如下:

@echo off
:loop
echo 启动处理脚本...
python your_script_name.py
echo 脚本退出,检查处理状态...
if exist progress.json (
    echo 未处理完成,重启脚本...
    goto loop
) else (
    echo 所有文件处理完成,退出。
)
pause

逻辑:脚本退出后检查progress.json是否存在,存在则说明未处理完,自动重启;不存在则处理完成,退出循环。

4. 高负载处理优化建议

  • 多进程并行:利用CPU多核并行处理不同文件(需加锁避免进度文件冲突)
    from multiprocessing import Pool, Lock
    lock = Lock()
    
    def process_file_with_lock(file):
        with lock:
            progress = load_progress()
            if file in progress["completed_files"]:
                return
        process_single_file(file)
    
    # 在main函数中替换剩余文件处理循环:
    remaining_files = [f for f in all_files if f not in completed_files]
    with Pool(processes=4) as pool:  # 根据CPU核心数调整进程数
        pool.map(process_file_with_lock, remaining_files)
    
  • 日志持久化:将日志写入文件,方便排查中断原因
    logging.basicConfig(
        level=logging.INFO,
        format='%(asctime)s - %(levelname)s - %(message)s',
        handlers=[logging.FileHandler("processing.log"), logging.StreamHandler()]
    )
    
  • 细粒度异常捕获:避免单条记录处理失败导致整个文件中断
    def clean_process(row):
        try:
            # 你的NLP清洗逻辑
            return row
        except Exception as e:
            logging.error(f"行处理失败: {str(e)},跳过该行")
            return None
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 13:52:48