Python多进程下向继承Process的类传递Task代理对象失效问题
问题根本原因
你遇到的问题和Task代理的实现无关,核心是多进程内存空间隔离的特性,以及对multiprocessing.Process的行为误解:
multiprocessing.Process的__init__方法是在父进程中执行的,只有run方法会在子进程的独立内存空间中运行。你在Scheduler中调用的worker.change_task(task)完全是在父进程地址空间执行,修改的是父进程中Worker实例的current_task属性,和子进程中的Worker实例没有任何关联。- 子进程启动后,内存空间和父进程完全独立,父进程对自己持有的Worker实例属性的修改不会同步到子进程,因此子进程中
current_task永远是初始化时的None。
解决方案
最简洁可靠的方案是通过多进程安全的通信管道传递Task代理,推荐用multiprocessing.Queue实现,修改代码如下:
import multiprocessing from multiprocessing import Queue class Worker(multiprocessing.Process): def __init__(self): super().__init__() # 跨进程队列,用于父进程向子进程传递Task代理 self._task_queue = Queue() self.current_task = None self.start() def change_task(self, new_task:Task): # 父进程调用该方法时,将Task代理塞入队列即可 self._task_queue.put(new_task) def run(self): while True: # 子进程阻塞等待接收新任务 self.current_task = self._task_queue.get() self.current_task.state = 'accepted-idle' self.current_task.state = 'accepted-running' # 执行任务逻辑 self.current_task.state = 'finished'
如果需要支持任务抢占、中途切换任务的场景,可以把队列读取改成非阻塞模式,定期检测是否有新任务传入:
def run(self): while True: # 非阻塞检测是否有新任务 try: new_task = self._task_queue.get(block=False) self.current_task = new_task self.current_task.state = 'accepted-idle' except multiprocessing.queues.Empty: pass if self.current_task is not None: self.current_task.state = 'accepted-running' # 执行任务逻辑 self.current_task.state = 'finished'
其他可选方案
如果不想用队列,也可以选择:
- 将
current_task属性也封装为Manager管理的代理对象,同时配套读写锁避免竞态问题 - 改用
multiprocessing.Manager.Namespace存储Worker的共享状态,所有跨进程修改的属性都放到这个共享命名空间中
内容的提问来源于stack exchange,提问作者Nxt-1
相关产品推荐
相关产品推荐

