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

Python多进程中如何将Master类状态传递给Worker进程

解决Python多进程中Master状态共享问题

问题根源

Python多进程采用**复制写时(Copy-on-Write)**机制,当你把Master实例传给子进程时,子进程拿到的是主进程中Master的独立副本,而非引用。Master的calculate方法在单独进程中更新的是原实例的状态,子进程里的副本完全不受影响,自然拿不到更新后的状态。

方案一:使用multiprocessing.Manager共享状态

这是最简洁的实现方式,Manager可以创建跨进程共享的对象,直接替代Master里的普通字典和锁:

修改后的Master类

import multiprocessing
import time

class Master:
    def __init__(self, manager):
        self.number = 0
        # 用Manager创建跨进程共享的字典和锁
        self.state = manager.dict()
        self.lock = manager.Lock()

    def calculate(self):
        while True:
            self.number += 1
            with self.lock:
                self.state["data"] = self.number
            time.sleep(1)

    def get_state(self):
        with self.lock:
            # 返回共享字典的副本,避免外部直接修改共享对象
            return dict(self.state)

修改后的主进程代码

def worker_process(worker, queue):
    while True:
        if not queue.empty():
            data = queue.get()
            worker.process(data, queue.qsize())

def main():
    # 初始化Manager,自动管理资源释放
    with multiprocessing.Manager() as manager:
        master = Master(manager)
        workers = [Worker(master, i) for i in range(multiprocessing.cpu_count())]

        queue = multiprocessing.Queue()

        # 启动Worker进程
        for worker in workers:
            multiprocessing.Process(target=worker_process, args=(worker, queue)).start()

        # 启动Master计算进程
        master_process = multiprocessing.Process(target=master.calculate)
        master_process.start()

        # 启动Eel进程
        eel_process = multiprocessing.Process(target=start_eel, args=(queue,))
        eel_process.start()

        # 等待进程结束(按需选择)
        master_process.join()
        eel_process.join()

if __name__ == "__main__":
    main()

Worker类无需修改,现在它拿到的Master实例里的state是跨进程共享的,调用get_state()就能获取最新状态。

方案二:用请求-响应队列实现状态查询(对应你提到的队列思路)

如果不想依赖Manager,可以用两个队列实现Worker与Master的请求-响应交互,通过Worker ID匹配对应结果:

实现细节

  1. 创建请求队列(Worker → Master)和响应队列(Master → Worker)
  2. Master进程同时运行两个任务:更新状态 + 监听请求队列并返回对应结果
  3. Worker查询状态时,发送带自身ID的请求,再从响应队列中筛选属于自己的结果

修改后的代码示例

Master类新增请求处理逻辑

class Master:
    def __init__(self):
        self.number = 0
        self.state = {}
        self.lock = multiprocessing.Lock()

    def calculate(self):
        while True:
            self.number += 1
            with self.lock:
                self.state = {"data": self.number}
            time.sleep(1)

    def get_state(self):
        with self.lock:
            return self.state.copy()
    
    def handle_requests(self, request_queue, response_queue):
        while True:
            if not request_queue.empty():
                worker_id = request_queue.get()
                state = self.get_state()
                response_queue.put((worker_id, state))

Worker类修改状态查询逻辑

class Worker:
    def __init__(self, cpu_id, request_queue, response_queue):
        self.cpu_id = cpu_id
        self.request_queue = request_queue
        self.response_queue = response_queue
        self.busy = False

    def get_state(self):
        # 发送带自身ID的查询请求
        self.request_queue.put(self.cpu_id)
        # 循环等待匹配自身ID的响应
        while True:
            resp_worker_id, state = self.response_queue.get()
            if resp_worker_id == self.cpu_id:
                return state

    def process(self, data, queue_size):
        print(f"Queue size: {queue_size}")
        self.busy = True
        state = self.get_state()
        print(f"Processing {[x - self.cpu_id for x in data]} using state {state} on CPU {self.cpu_id}")
        time.sleep(2)
        self.busy = False

主进程代码调整

def worker_process(worker, queue):
    while True:
        if not queue.empty():
            data = queue.get()
            worker.process(data, queue.qsize())

def main():
    master = Master()
    # 创建请求和响应队列
    request_queue = multiprocessing.Queue()
    response_queue = multiprocessing.Queue()

    # 初始化Worker时传入队列和ID,不再传递Master实例
    workers = [Worker(i, request_queue, response_queue) for i in range(multiprocessing.cpu_count())]

    task_queue = multiprocessing.Queue()

    # 启动Worker进程
    for worker in workers:
        multiprocessing.Process(target=worker_process, args=(worker, task_queue)).start()

    # 启动Master的两个进程:计算状态 + 处理请求
    master_calc_process = multiprocessing.Process(target=master.calculate)
    master_req_process = multiprocessing.Process(target=master.handle_requests, args=(request_queue, response_queue))
    master_calc_process.start()
    master_req_process.start()

    eel_process = multiprocessing.Process(target=start_eel, args=(task_queue,))
    eel_process.start()

    master_calc_process.join()
    master_req_process.join()
    eel_process.join()

if __name__ == "__main__":
    main()

方案对比

  • Manager方案:代码改动极小,实现简单,适合大多数场景。缺点是基于RPC通信,性能略低于直接用队列。
  • 请求-响应队列方案:完全手动控制进程间通信,性能更高,但代码复杂度上升,需要处理请求匹配逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 18:35:14