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匹配对应结果:
实现细节
- 创建请求队列(Worker → Master)和响应队列(Master → Worker)
- Master进程同时运行两个任务:更新状态 + 监听请求队列并返回对应结果
- 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
相关产品推荐
相关产品推荐

