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
相关产品推荐
相关产品推荐

