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

如何按顺序使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 19:09:00