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

Python multiprocessing如何最优启动存在依赖关系的多进程?

Python multiprocessing带依赖的进程启动最优实现

实现这种带前置依赖的进程启动逻辑,最可靠的方案是使用Python标准库multiprocessing提供的Event同步原语,无需固定时长sleep,也不需要轮询队列空状态,完全由进程主动触发启动信号。

核心实现逻辑

  • 创建两个Event实例分别作为两个下游进程的启动触发信号
  • 上游进程完成第一次数据写入后,主动触发对应事件
  • 主进程阻塞等待事件触发后,再启动对应的下游进程

同时原有代码中存在一个隐藏问题:消费者进程用while not queue.empty()作为循环判断条件,一旦生产者生产速度慢于消费速度,队列临时为空时消费者会直接退出,后续生产的新数据无法被处理。可以通过哨兵值(比如None)标识生产结束,避免这个问题。

优化后完整代码

import multiprocessing
import keyboard
import time

def getData(queue_raw, raw_ready_event):
    for num in range(1000):
        queue_raw.put(num)
        print(f"getData: put {num} in queue_raw")
        # 第一次写入数据后触发启动信号,只触发一次
        if num == 0:
            raw_ready_event.set()
    # 生产结束写入哨兵值
    queue_raw.put(None)
    while True:
        if keyboard.read_key() == "s":
            break

def calcFeatures(queue_raw, queue_features, feature_ready_event):
    while True:
        data = queue_raw.get()
        # 收到哨兵值退出循环
        if data is None:
            queue_features.put(None)
            break
        res = data ** 2
        queue_features.put(res)
        print(f"calcFeatures: put {res} in queue_features")
        # 第一次写入特征后触发启动信号
        if res == 0:
            feature_ready_event.set()

def sendFeatures(queue_features):
    while True:
        feature = queue_features.get()
        if feature is None:
            break
        print(f"sendFeatures: put {feature} out")

if __name__ == "__main__":
    queue_raw = multiprocessing.Queue()
    queue_features = multiprocessing.Queue()
    # 定义两个启动触发事件
    raw_ready = multiprocessing.Event()
    feature_ready = multiprocessing.Event()

    processes = [
        multiprocessing.Process(target=getData, args=(queue_raw, raw_ready)),
        multiprocessing.Process(target=calcFeatures, args=(queue_raw, queue_features, feature_ready)),
        multiprocessing.Process(target=sendFeatures, args=(queue_features,))
    ]

    # 启动第一个进程
    processes[0].start()
    # 阻塞等待第一个进程写入第一条数据
    raw_ready.wait()
    # 启动第二个进程
    processes[1].start()
    # 阻塞等待第二个进程写入第一条特征
    feature_ready.wait()
    # 启动第三个进程
    processes[2].start()

    for p in processes:
        p.join()

方案说明

  • Event.wait()会一直阻塞直到事件被set(),不会占用额外CPU资源,比轮询队列空状态效率高很多
  • 事件触发是一次性的,只要触发过一次wait()就会直接返回,不会重复阻塞
  • 新增的哨兵值逻辑可以保证所有数据都被处理完,消费者不会提前退出

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 03:36:03