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

如何给递归函数内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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 23:54:06