Python multiprocessing多文件夹并行处理及子进程状态追踪方案咨询
多进程改造实现方案
核心实现逻辑
- 废弃全局变量
count,将每个临时文件夹的序号、对应war包路径、metadata路径作为独立参数传给子任务,避免多进程间变量冲突 - 使用
multiprocessing.Pool控制最大并行数,可根据服务器性能和Nexus承载能力调整并行数量 - 每个子任务内部捕获所有异常、命令执行返回值,执行成功返回
True,失败返回False及错误信息 - 主进程收集所有子任务的执行结果,只要存在任意一个失败的任务,就标记整体构建失败
改造后的完整代码
import re, shutil, os import multiprocessing from pathlib import Path target = "/home/work" file_path = target + "/file.txt" # 不要用dict作为变量名,会覆盖Python内置类 war_metadata_map = {} # 最大并行数,可根据实际情况调整为100 MAX_WORKERS = 50 def process_single_folder(task_params): """单文件夹处理任务,每个子进程独立执行该函数""" war_path, metadata_path, task_id = task_params tmp_folder_path = f"{target}/tmp{task_id}" try: # 1. 创建临时文件夹,拷贝对应资源 os.mkdir(tmp_folder_path) shutil.copy(war_path, tmp_folder_path) shutil.copy(metadata_path, tmp_folder_path) # 2. 执行处理命令,所有命令返回值非0都判定为执行失败 # 替换为你实际的构建标签命令 ret = os.system("<1st command to build the label>") if ret != 0: raise Exception(f"构建标签失败,返回码:{ret}") ret = os.system("<2nd command to build the package>") if ret != 0: raise Exception(f"构建包失败,返回码:{ret}") # 这里替换为你原来的文件操作逻辑,filetype固定为war可直接使用,也可作为参数传入 # <multiple file manipulations> ret = os.system("<curl command to upload it to the Nexus>") if ret != 0: raise Exception(f"上传Nexus失败,返回码:{ret}") # 执行成功返回状态、任务ID、空错误信息 return (True, task_id, None) except Exception as e: # 执行失败返回状态、任务ID、错误详情 return (False, task_id, str(e)) if __name__ == "__main__": # 读取文件组装war包和metadata的映射关系 with open(file_path, 'r') as f: lines = f.read().splitlines() for i, line in enumerate(lines): match = re.match(r".*.war", line) if match: j = i-1 if i > 1 else 0 for k in range(j, i): war_metadata_map[match.string] = lines[k] # 组装任务参数列表 task_list = [] for idx, (war_path, metadata_path) in enumerate(war_metadata_map.items(), start=1): task_list.append((war_path, metadata_path, idx)) # 初始化进程池,maxtasksperchild=1可避免子进程内存泄漏 pool = multiprocessing.Pool(processes=MAX_WORKERS, maxtasksperchild=1) # 批量提交所有任务,异步执行 results = pool.map(process_single_folder, task_list) # 关闭进程池,等待所有子任务执行完成 pool.close() pool.join() # 检查所有任务执行结果 all_success = True failed_tasks = [] for res in results: success, task_id, err_msg = res if not success: all_success = False failed_tasks.append(f"任务{task_id}:{err_msg}") if all_success: print("所有任务执行成功,整体构建成功") else: print("整体构建失败,失败的任务如下:") for err in failed_tasks: print(f"- {err}") # 非0退出码标记整体失败 exit(1)
注意事项
- 并行数不建议直接设置为100,可先从20~50开始测试,避免Nexus上传带宽被打满、服务器IO过载
- 建议将
os.system替换为subprocess.run,可以更精准的捕获命令的输出和返回值,方便排查失败原因 - 可按需添加临时文件夹执行完自动清理的逻辑,避免占用磁盘空间
- 22GB内存的容器完全可以支撑100个并行任务,瓶颈大概率出现在网络上传、磁盘IO环节
内容的提问来源于stack exchange,提问作者Chel MS
相关产品推荐
相关产品推荐

