如何给递归函数内os.fork()创建的子进程设置并发数量上限
递归fork进程并发数上限实现方案
核心思路是使用跨进程共享的信号量控制并发数,实现简单且性能稳定,具体操作如下:
实现步骤
- 提前初始化一个进程间共享的信号量,初始值设为你需要的并发上限30
- 每次执行
os.fork()前先获取信号量,获取失败时自动阻塞,直到有运行中的子进程退出释放资源 - 子进程执行完递归任务退出前,必须释放信号量,把名额让给等待中的新进程
修复后可运行的代码示例
首先你需要在代码最开头初始化信号量,注意不要放在递归函数内部:
import os import re from multiprocessing import Semaphore # 初始化全局共享信号量,并发上限设为30 MAX_PROCESS = 30 sem = Semaphore(MAX_PROCESS) pathpattern = r'你的路径匹配正则' # 替换为你实际的匹配规则 def recursive_copying(file, target_path): newlines = [] processes = [] with open(file, 'r', encoding='utf-8') as f: for line in f: modified_line = line # 替换为你实际的行处理逻辑 # 匹配到路径时创建子进程处理 if re.match(pathpattern, line): file_line = line.strip() # 注意处理换行符 # 先获取信号量,占一个进程名额 sem.acquire() pid = os.fork() if pid == 0: # 子进程逻辑,执行完必须释放信号量 try: recursive_copying(file_line, target_path) finally: sem.release() os._exit(0) else: processes.append(pid) # 处理当前行内容 newlines.append(modified_line) # 写入修改后的新文件 new_file = os.path.join(target_path, os.path.basename(file)) # 替换原来的路径拼接逻辑,更安全 with open(new_file, 'w+', encoding='utf-8') as f: f.writelines(newlines) # 等待当前函数创建的所有子进程退出 for pid in processes: os.waitpid(pid, 0) return
注意事项
- 信号量必须在所有fork操作之前初始化,放在全局作用域即可,不要放在递归函数内部,否则每个进程都会生成独立的信号量,无法实现全局计数
- 子进程的信号量释放逻辑放在
try-finally块中,避免子进程运行报错时无法释放信号量导致死锁 - 原代码中的路径拼接、换行符处理、行修改逻辑需要根据你自己的业务调整,示例中仅做了基础补全
内容的提问来源于stack exchange,提问作者Regretful
相关产品推荐
相关产品推荐

