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
相关产品推荐
相关产品推荐

