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

如何为multiprocessing.Pool分配两个函数实现并行执行?

用multiprocessing.Pool同时运行initiate和formula的方案

咱们先理清楚:multiprocessing.Pool的设计初衷是批量处理同类型并行任务,但这不代表它不能同时跑两个不同的函数——只要处理好进程间的数据共享和任务提交逻辑,刚好能解决你之前速度慢的问题(大概率是之前只开了一个formula进程,Pool可以帮你把多核CPU用满)。

核心思路

你的场景本质是生产者-消费者模式:initiate是生产者(持续从WebSocket拉取数据),formula是消费者(基于数据做计算和场景判断)。用Pool的话,咱们可以:

  1. 用跨进程共享的数据结构传递数据(多进程内存不共享,普通列表没法直接用)
  2. 给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()

关键注意事项

  1. 数据共享性能优化:如果你的数据量很大,Manager.List的性能可能不足,建议换成multiprocessing.Queue——它是专门为进程间通信设计的,生产者用put()存数据,消费者用get()取数据,线程/进程安全且速度更快。只需要把shared_data = manager.List()换成shared_data = multiprocessing.Queue(),然后修改对应的数据操作逻辑即可。

  2. Pool的任务分配优势:Pool会自动把任务分配给空闲进程,不用手动管理,比手动创建多个Process更省心,还能避免资源浪费。

  3. 避免CPU空转:在formula里加time.sleep(0.05)很重要,不然进程会一直循环检查数据,把CPU占满。如果用Queue,get()方法默认是阻塞式的,没有数据时进程会自动休眠,不用额外加休眠。

  4. 优雅终止:代码中加入了KeyboardInterrupt捕获,按Ctrl+C时,进程能正常关闭WebSocket并退出。

为什么比之前的Process快?

之前你用multiprocessing.Process同时跑1个initiate和1个formula,只能用到2个CPU核心;现在用Pool可以把剩下的所有核心都用来跑formula,相当于多开了N-1个消费者进程,自然能大幅提升处理速度。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:09:05