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

Python Multiprocessing:交替多进程与单进程执行问题求助

多进程交替采集-处理的同步解决方案

问题分析

你的核心需求是实现主进程与N个子进程的严格交替执行:子进程批量采集数据→主进程统一处理→子进程再采集,循环往复。原方案存在两个关键问题:

  • 同步逻辑漏洞:使用单个Semaphore时,子进程释放信号量后会立即重新抢占,导致主进程无法一次性获取所有N个信号量,无法进入处理阶段。
  • 数据共享错误:子进程的buffer与主进程的Generator实例内存空间隔离,主进程调用get_values()无法获取子进程采集的数据。

推荐方案:使用Barrier同步+Queue传递数据

针对你的需求,multiprocessing.Barrier是最贴合的同步工具——它可以让指定数量的进程在某个点同步,完美适配"所有子进程完成采集后主进程再处理,主进程处理完成后子进程再继续"的交替逻辑。同时用Queue实现子进程到主进程的安全数据传递。

修正后的代码

import multiprocessing as mp
import time

def generator_task(process_id: int, k: int, collect_barrier: mp.Barrier, process_barrier: mp.Barrier, queue: mp.Queue):
    while True:
        # 阶段1:采集K个值
        buffer = []
        for i in range(k):
            buffer.append(f"proc_{process_id}_val_{i}")
            time.sleep(0.3)  # 模拟采集耗时
        print(f"Process {process_id} finished collecting: {buffer}")
        
        # 将数据发送给主进程
        queue.put((process_id, buffer))
        
        # 等待所有子进程完成采集,主进程进入处理阶段
        collect_barrier.wait()
        
        # 等待主进程完成处理,进入下一轮采集
        process_barrier.wait()

if __name__ == '__main__':
    n_proc = 4
    k_per_proc = 2
    
    # 初始化同步工具:两个Barrier,参与方为子进程数+主进程
    collect_barrier = mp.Barrier(n_proc + 1)
    process_barrier = mp.Barrier(n_proc + 1)
    
    # 初始化队列用于传递采集的数据
    data_queue = mp.Queue()
    
    # 启动子进程
    processes = []
    for proc_id in range(n_proc):
        proc = mp.Process(
            target=generator_task,
            args=(proc_id, k_per_proc, collect_barrier, process_barrier, data_queue)
        )
        processes.append(proc)
        proc.start()
    
    print("All child processes started")
    
    try:
        while True:
            # 阶段2:等待所有子进程完成采集
            print("\nWaiting for all processes to finish collecting...")
            collect_barrier.wait()
            
            # 从队列收集所有子进程的数据
            all_values = []
            for _ in range(n_proc):
                proc_id, vals = data_queue.get()
                all_values.extend(vals)
            
            print(f"Collected all data: {all_values}")
            # 模拟主进程处理逻辑
            print("Main process processing data...")
            time.sleep(1)
            print("Main process finished processing\n")
            
            # 通知所有子进程可以开始下一轮采集
            process_barrier.wait()
    except KeyboardInterrupt:
        # 清理子进程
        for proc in processes:
            proc.terminate()
        print("\nProcesses terminated")

方案说明

  1. 同步逻辑:
    • collect_barrier:所有子进程完成采集后调用wait(),主进程也调用wait();当所有n_proc+1个进程都到达同步点时,主进程开始处理数据,子进程进入等待状态。
    • process_barrier:主进程完成处理后调用wait(),所有子进程也调用wait();同步完成后子进程进入下一轮采集,主进程回到等待状态。
  2. 数据传递:使用Queue安全地在子进程和主进程之间传递数据,彻底解决多进程内存隔离导致的数据无法共享问题。
  3. 优势:逻辑清晰,完全贴合交替执行的需求,无需手动管理信号量抢占,Barrier会自动处理同步点。

替代方案:使用Event+共享计数器

如果不想用Barrier,也可以用两个Event配合一个共享计数器实现:

  • collect_done:主进程等待该事件,当所有子进程完成采集后触发。
  • process_done:子进程等待该事件,主进程处理完成后触发。
  • 共享计数器:记录完成采集的子进程数量,达到n_proc时触发collect_done。

不过Barrier方案更简洁,无需手动维护计数器状态。

内容的提问来源于stack exchange,提问作者Federico Taschin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 15:51:39