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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 11:09:03