如何为multiprocessing.Pool分配两个函数实现并行执行?
咱们先理清楚:multiprocessing.Pool的设计初衷是批量处理同类型并行任务,但这不代表它不能同时跑两个不同的函数——只要处理好进程间的数据共享和任务提交逻辑,刚好能解决你之前速度慢的问题(大概率是之前只开了一个formula进程,Pool可以帮你把多核CPU用满)。
核心思路
你的场景本质是生产者-消费者模式:initiate是生产者(持续从WebSocket拉取数据),formula是消费者(基于数据做计算和场景判断)。用Pool的话,咱们可以:
- 用跨进程共享的数据结构传递数据(多进程内存不共享,普通列表没法直接用)
- 给Pool提交1个
initiate任务,再提交多个formula任务(用剩下的CPU核心)
具体实现步骤和代码示例
第一步:处理跨进程数据共享
普通列表在多进程中是各自独立的,必须用multiprocessing.Manager提供的共享结构,比如Manager.List或者Queue(后者更适合生产者-消费者场景,性能更好)。这里先以Manager.List为例,后面会补充Queue的优化方案。
第二步:编写可在Pool中运行的函数
要把共享数据结构作为参数传给函数,让生产者和消费者能共用同一份数据:
import multiprocessing from multiprocessing import Pool, Manager import time # 替换成你实际使用的WebSocket库和处理逻辑 import websocket def process_websocket_data(raw_data): # 模拟你对WebSocket原始数据的处理逻辑 return {"value": float(raw_data), "timestamp": time.time()} def check_scenario(data, scenario_id): # 模拟你的场景判断逻辑,不同scenario_id对应不同规则 if scenario_id == 0: return data["value"] > 100 elif scenario_id == 1: return data["value"] < 50 return False def perform_action(data): # 模拟你要执行的操作 print(f"触发操作:{data}") # 生产者函数:WebSocket拉取数据并存入共享列表 def initiate(shared_data): ws = websocket.WebSocket() ws.connect("ws://your-websocket-endpoint") try: while True: raw_data = ws.recv() processed_data = process_websocket_data(raw_data) shared_data.append(processed_data) time.sleep(0.1) # 模拟数据接收间隔,根据实际情况调整 except KeyboardInterrupt: ws.close() # 消费者函数:读取共享数据,符合场景则执行操作 def formula(shared_data, scenario_id): try: while True: if len(shared_data) > 0: # 取出队列头部数据,避免重复处理 data = shared_data.pop(0) if check_scenario(data, scenario_id): perform_action(data) # 加小休眠避免空转占用过多CPU time.sleep(0.05) except KeyboardInterrupt: pass
第三步:用Pool启动任务
在主进程中创建Pool,提交生产者和多个消费者任务:
if __name__ == "__main__": # 创建共享数据容器 manager = Manager() shared_data = manager.List() # 获取CPU核心数,留1个给生产者,剩下的全分配给消费者 cpu_count = multiprocessing.cpu_count() consumer_count = cpu_count - 1 # 用with语句自动管理Pool生命周期,避免资源泄漏 with Pool(cpu_count) as pool: # 异步提交生产者任务 pool.apply_async(initiate, args=(shared_data,)) # 提交多个消费者任务,每个对应不同场景 for scenario_id in range(consumer_count): pool.apply_async(formula, args=(shared_data, scenario_id)) # 等待所有任务完成(因是无限循环,会一直运行直到按Ctrl+C终止) pool.close() pool.join()
关键注意事项
数据共享性能优化:如果你的数据量很大,
Manager.List的性能可能不足,建议换成multiprocessing.Queue——它是专门为进程间通信设计的,生产者用put()存数据,消费者用get()取数据,线程/进程安全且速度更快。只需要把shared_data = manager.List()换成shared_data = multiprocessing.Queue(),然后修改对应的数据操作逻辑即可。Pool的任务分配优势:Pool会自动把任务分配给空闲进程,不用手动管理,比手动创建多个
Process更省心,还能避免资源浪费。避免CPU空转:在
formula里加time.sleep(0.05)很重要,不然进程会一直循环检查数据,把CPU占满。如果用Queue,get()方法默认是阻塞式的,没有数据时进程会自动休眠,不用额外加休眠。优雅终止:代码中加入了
KeyboardInterrupt捕获,按Ctrl+C时,进程能正常关闭WebSocket并退出。
为什么比之前的Process快?
之前你用multiprocessing.Process同时跑1个initiate和1个formula,只能用到2个CPU核心;现在用Pool可以把剩下的所有核心都用来跑formula,相当于多开了N-1个消费者进程,自然能大幅提升处理速度。
内容的提问来源于stack exchange,提问作者Thomas Bernhard

