如何按顺序使用Python multiprocessing避免文件并发访问冲突
问题描述
现有Python脚本包含两个执行环节:
- 环节1:加载并解压文件
- 环节2:处理解压完成的文件
未引入多进程时程序按顺序执行,先完成所有文件解压再启动处理逻辑。引入multiprocessing后,解压操作未完成时文件处理流程就已启动,大数据量(100+文件)场景下会出现并发文件访问错误,抛出PermissionError: [WinError 32] The process cannot access the file because it is being used by another process:;小数据量(约30个文件)场景下因解压速度快可正常运行。
实现需求
保留多进程的性能优势,仅在所有文件解压完成后再启动文件处理流程。
问题根因
Windows系统下multiprocessing创建子进程时会重新导入整个主模块,当前文件处理代码没有放在if __name__ == '__main__'判断块内,导致子进程启动时就会执行处理逻辑,出现解压和处理并行的冲突。另外pool.map本身是阻塞方法,会等待所有子进程的解压任务全部执行完成才会继续向下执行,刚好可以满足等所有解压完成再处理的需求。
修复方案
将所有文件处理逻辑全部挪到if __name__ == '__main__'块内,放在解压任务执行完成之后,同时补充进程池资源释放逻辑,另外修正原代码中filter_row函数的参数笔误(形参为r,内部错误使用row)。
修复后的完整代码如下:
import os import csv import collections import datetime import zipfile import re import shutil import fnmatch from pathlib import Path import ntpath import configparser from multiprocessing import Pool def generate_file_lists(): # 可修改为实际路径 data_files = 'c:\\desktop\\DataEnergy' pattern = '*.zip' last_root = None args = [] for root, dirs, files in os.walk(data_files): for filename in fnmatch.filter(files, pattern): if root != last_root: last_root = root if args: yield args args = [] args.append((root, filename)) if args: yield args def unzip(file_list): """ file_list 是同目录下的 (root, filename) 元组列表 """ # 可修改为实际路径 counter_part = 'c:\\desktop\\CounterPart' for root, filename in file_list: path = os.path.join(root, filename) date_zipped_file_s = re.search('-(.\d+)-', filename).group(1) date_zipped_file = datetime.datetime.strptime(date_zipped_file_s, '%Y%m%d').date() # 创建新目录路径 new_dir = os.path.normpath(os.path.join(os.path.relpath(path, start='c:\\desktop\\DataEnergy'), "..")) # 拼接输出根路径 new = os.path.join(counter_part, new_dir) # 创建目录 if not os.path.exists(new): os.makedirs(new) zipfile.ZipFile(path).extractall(new) # 获取解压后的所有文件 files = os.listdir(new) # 重命名文件 for file in files: filesplit = os.path.splitext(os.path.basename(file)) if not re.search(r'_\d{8}.', file): os.rename(os.path.join(new, file), os.path.join(new, filesplit[0]+'_'+date_zipped_file_s+filesplit[1])) # 仅读取符合条件的行 def filter_row(r, missing_date): if set(r).intersection({'Audi', 'Mercedes', 'Volkswagen'}): if len(r) > 24 and r[24].isdigit(): date_time = datetime.datetime.fromisoformat(r[0]) if len(r) > 3 else True condition_3 = date_time.date() == missing_date if len(r) > 3 else True return condition_3 return False # Windows系统必须加该判断 if __name__ == '__main__': # 执行解压任务 pool = Pool(13) pool.map(unzip, generate_file_lists()) # 释放进程池资源 pool.close() pool.join() print('所有文件解压完成!') # 启动处理流程 all_missing_dates = ['20210701', '20210702'] missing_dates = [datetime.datetime.strptime(i, "%Y%m%d").date() for i in all_missing_dates] dates_to_process = [] root = Path('.\\middle_stage').resolve() print("开始读取数据") data_per_date = dict() for missing_date in missing_dates: print("\t正在读取日期数据: ", missing_date) files=[fn for fn in (e for e in root.glob(f"**/*_{missing_date:%Y%m%d}.txt") if e.is_file())] if len(files) != 13: continue dates_to_process.append(missing_date) vehicle_loc_dict = collections.defaultdict(list) for file in files: with open(file, 'r', encoding='utf-8') as log_file: reader = csv.reader(log_file, delimiter = ',') next(reader) # 跳过表头 for row in reader: if filter_row(row, missing_date): print('filter_row 执行完成!') data_per_date[missing_date] = vehicle_loc_dict
内容的提问来源于stack exchange,提问作者Mediterráneo
相关产品推荐
相关产品推荐

