使用ProcessPoolExecutor.map处理文件时出现queue.FULL错误求助
解决ProcessPoolExecutor.map()出现queue.FULL错误的方法
嘿,我之前处理大量文件并行任务时也碰到过这个问题!当你用ProcessPoolExecutor.map()一次性提交400个文件处理任务时,很容易因为任务提交速度远快于进程处理速度,导致内部任务队列被塞满,触发queue.FULL错误。下面是几个实用的解决办法:
1. 调整chunksize参数(最简单高效)
ProcessPoolExecutor.map()默认会把每个任务单独放入队列,当任务数量远大于工作进程数时,队列很容易被撑爆。通过设置chunksize,可以让每个工作进程一次性获取一批任务,减少队列的任务堆积。
你可以根据工作进程数来计算合适的chunksize,比如:
from concurrent.futures import ProcessPoolExecutor import glob import os # 你的文件处理函数 def excel_to_csv(infile, outfile): # 这里替换成你的Excel转CSV逻辑,比如用pandas: # import pandas as pd # df = pd.read_excel(infile) # df.to_csv(outfile, index=False) pass input_path = "你的输入路径" infiles = glob.glob(os.path.join(input_path, '**/*.xls'), recursive=True) + glob.glob(os.path.join(input_path, '**/*.xlsx'), recursive=True) outfiles = [os.path.join(os.path.dirname(f), f"{os.path.basename(f).split('.')[0]}.csv") for f in infiles] max_workers = 4 # 可根据CPU核心数调整,默认是CPU核心数 chunksize = len(infiles) // max_workers + 1 # 平均分配任务,有余数就+1 with ProcessPoolExecutor(max_workers=max_workers) as executor: # 传入chunksize参数 executor.map(excel_to_csv, infiles, outfiles, chunksize=chunksize)
2. 手动控制任务提交,避免一次性塞满队列
如果chunksize的方式还是没解决问题,你可以用submit()配合as_completed(),手动控制每次提交的任务数量,确保队列不会过载。这种方式更灵活,还能实时处理任务结果或异常:
from concurrent.futures import ProcessPoolExecutor, as_completed input_path = "你的输入路径" infiles = glob.glob(os.path.join(input_path, '**/*.xls'), recursive=True) + glob.glob(os.path.join(input_path, '**/*.xlsx'), recursive=True) outfiles = [os.path.join(os.path.dirname(f), f"{os.path.basename(f).split('.')[0]}.csv") for f in infiles] max_workers = 4 task_pairs = list(zip(infiles, outfiles)) with ProcessPoolExecutor(max_workers=max_workers) as executor: # 先提交第一批任务,数量等于工作进程数 futures = { executor.submit(excel_to_csv, infile, outfile): (infile, outfile) for infile, outfile in task_pairs[:max_workers] } remaining_tasks = task_pairs[max_workers:] while futures: # 遍历已完成的任务 for future in as_completed(futures): task_info = futures[future] try: # 获取任务结果(如果有返回值的话) future.result() print(f"成功处理文件: {task_info[0]}") except Exception as e: print(f"处理文件 {task_info[0]} 出错: {str(e)}") # 如果还有剩余任务,提交下一个 if remaining_tasks: next_infile, next_outfile = remaining_tasks.pop(0) new_future = executor.submit(excel_to_csv, next_infile, next_outfile) futures[new_future] = (next_infile, next_outfile) # 移除已处理完成的任务 del futures[future]
3. 优化处理函数的执行效率
队列满的核心问题是任务处理速度跟不上提交速度,所以优化你的Excel转CSV函数也很关键:
- 使用更高效的库:比如用
pandas替代手动解析Excel,或者指定更快的引擎(read_excel(engine='openpyxl')针对xlsx,engine='xlrd'针对xls); - 减少不必要的内存操作:比如读取Excel时只加载需要的列,避免全量加载;
- 尽量避免在处理函数中进行IO密集型操作(比如频繁读写文件),如果必须做,尽量批量处理。
内容的提问来源于stack exchange,提问作者Manuel Pasieka
相关产品推荐
相关产品推荐

