Python线程/进程中运行任务的对象更新不生效问题求助
我开发了一个主应用程序,会创建任务交由独立的thread/process(两种方案均已尝试)执行。具体实现为实例化Handler对象并传入任务。后续主程序获得需添加至运行中任务的额外信息,通过调用Handler类的update_task方法传递该信息。在update_task方法内部可看到更新后的任务,但Handler类的其他位置(如run方法)仍引用初始对象。
以下是简化后的代码,将两个类放在同一模块中:
#!/usr/bin/env python3 # encoding: utf-8 from pprint import pformat from multiprocessing import Process from threading import Event, Lock from copy import deepcopy # Handles a single task class Handler(Process): def __init__(self, task): super(Handler, self).__init__() self.__evt = Event() self.__task = deepcopy(task) self.__lock = Lock() print(f"Using Task: {task['id']}:\n {(task)}") print(f"{id(self)} ==> Initial Task ID: {id(self.__task)}") def update_task(self, task): print("*"*50) print(f"Updating the Task to {(task)}") with self.__lock: del self.__task self.__task = deepcopy(task) print(f"{id(self)} ==> New Task ID: {id(self.__task)}") print(f"The NEW TASK: {self.__task}") print("*"*50) def run(self): super(Handler, self).run() print("Starting the Main Thread") while not self.__evt.isSet(): with self.__lock: print(f"Task ID: {id(self.__task)} ==> {self.__task}") self.__evt.wait(0.5) def stop(self): print(f"Stopping Thread") super(Handler, self).terminate() self.__evt.set() # The main application class Driver: def __init__(self): self.__evt = Event() self.__task = { 'id': "Lonely Task", 'name': 'Testing Item', 'sources': [] } self.__task['sources'].append(self.__gen_entry()) print(f"The Task: {pformat(self.__task)}") handler = Handler(self.__task) handler.start() self.__evt.wait(1) print("Adding a new entry") self.__task['sources'].append(self.__gen_entry()) self.__evt.wait(1) print("*********** Now Updating the Handler ***************") handler.update_task(self.__task) self.__evt.wait(3) print("Now Stopping the Handler") handler.stop() print("Done!!") def __gen_entry(self): print("Generating entry") sz = len(self.__task['sources']) n = sz + 1 return { 'id': f"entry-{n}", 'rx': f"Rx{n}", 'fx': f"Fx{n}", } if __name__ == '__main__': Driver()
我尝试使用Thread和Process,结果一致;也尝试过使用/不使用copy/deepcopy、删除初始对象、使用/不使用Lock对象,但均未解决问题。我预期调用update_task更新任务后,run方法后续的打印信息会体现这些变化,但实际并未实现。以下是示例输出:
./main_driver.py Generating entry The Task: {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'fx': 'Fx1', 'id': 'entry-1', 'rx': 'Rx1'}]} Using Task: Lonely Task: {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} 139709461831752 ==> Initial Task ID: 139709454116256 Starting the Main Thread Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} Adding a new entry Generating entry Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} *********** Now Updating the Handler *************** ************************************************** Updating the Task to {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}, {'id': 'entry-2', 'rx': 'Rx2', 'fx': 'Fx2'}]} 139709461831752 ==> New Task ID: 139709454116328 The NEW TASK: {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}, {'id': 'entry-2', 'rx': 'Rx2', 'fx': 'Fx2'}]} ************************************************** Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} Task ID: 139709454116256 ==> {'id': 'Lonely Task', 'name': 'Testing Item', 'sources': [{'id': 'entry-1', 'rx': 'Rx1', 'fx': 'Fx1'}]} Now Stopping the Handler Stopping Thread Done!!
问题核心在于线程与进程的内存隔离机制差异,直接调用实例方法修改属性的逻辑在两种场景下不通用,需针对性调整:
1. 使用Thread(线程)场景
线程共享同一进程的内存空间,修改实例属性可直接被其他线程感知。只需调整Handler的继承类,并确保锁的正确使用:
from threading import Thread, Event, Lock from copy import deepcopy class Handler(Thread): def __init__(self, task): super().__init__() self.__evt = Event() self.__task = deepcopy(task) self.__lock = Lock() print(f"Using Task: {task['id']}:\n {(task)}") print(f"{id(self)} ==> Initial Task ID: {id(self.__task)}") def update_task(self, task): print("*"*50) print(f"Updating the Task to {(task)}") with self.__lock: self.__task = deepcopy(task) print(f"{id(self)} ==> New Task ID: {id(self.__task)}") print(f"The NEW TASK: {self.__task}") print("*"*50) def run(self): print("Starting the Main Thread") while not self.__evt.isSet(): with self.__lock: print(f"Task ID: {id(self.__task)} ==> {self.__task}") self.__evt.wait(0.5) def stop(self): print(f"Stopping Thread") self.__evt.set() self.join()
此时主线程调用update_task修改的是同一个实例的__task属性,run方法所在线程加锁后读取的会是最新引用,即可看到更新后的任务。
2. 使用Process(进程)场景
进程拥有完全独立的内存空间,主进程中修改Handler实例的属性对子进程无任何影响,必须使用**进程间通信(IPC)**机制传递更新:
方案1:使用Queue传递更新指令
在Handler中维护一个队列,主进程通过队列发送更新任务,子进程在run循环中主动检查并更新:
from multiprocessing import Process, Queue, Lock from threading import Event from copy import deepcopy class Handler(Process): def __init__(self, task): super().__init__() self.__evt = Event() self.__task = deepcopy(task) self.__lock = Lock() self.__update_queue = Queue() print(f"Using Task: {task['id']}:\n {(task)}") print(f"{id(self)} ==> Initial Task ID: {id(self.__task)}") def update_task(self, task): self.__update_queue.put(deepcopy(task)) print("*"*50) print(f"Sent update for task: {task}") print("*"*50) def run(self): print("Starting the Main Process") while not self.__evt.isSet(): # 优先处理队列中的更新 while not self.__update_queue.empty(): with self.__lock: self.__task = self.__update_queue.get() print(f"{id(self)} ==> Updated Task ID: {id(self.__task)}") # 输出当前任务状态 with self.__lock: print(f"Task ID: {id(self.__task)} ==> {self.__task}") self.__evt.wait(0.5) def stop(self): print(f"Stopping Process") self.__evt.set() self.join()
方案2:使用Manager创建共享对象
通过multiprocessing.Manager创建可跨进程共享的字典/列表,主进程修改共享对象时,子进程可直接感知变化:
from multiprocessing import Process, Manager, Lock from threading import Event from pprint import pformat class Handler(Process): def __init__(self, shared_task): super().__init__() self.__evt = Event() self.__task = shared_task self.__lock = Lock() print(f"Using Task: {shared_task['id']}:\n {(shared_task)}") print(f"{id(self)} ==> Initial Task ID: {id(self.__task)}") def update_task(self, task): print("*"*50) print(f"Updating the Task to {(task)}") with self.__lock: self.__task.clear() self.__task.update(task) print(f"{id(self)} ==> Task ID remains: {id(self.__task)}") print(f"The NEW TASK: {self.__task}") print("*"*50) def run(self): print("Starting the Main Process") while not self.__evt.isSet(): with self.__lock: print(f"Task ID: {id(self.__task)} ==> {self.__task}") self.__evt.wait(0.5) def stop(self): print(f"Stopping Process") self.__evt.set() self.join() # 修改Driver类,使用Manager创建共享对象 class Driver: def __init__(self): self.__evt = Event() with Manager() as manager: self.__task = manager.dict({ 'id': "Lonely Task", 'name': 'Testing Item', 'sources': manager.list() }) self.__task['sources'].append(self.__gen_entry()) print(f"The Task: {pformat(self.__task)}") handler = Handler(self.__task) handler.start() self.__evt.wait(1) print("Adding a new entry") self.__task['sources'].append(self.__gen_entry()) self.__evt.wait(1) print("*********** Now Updating the Handler ***************") handler.update_task(dict(self.__task)) self.__evt.wait(3) print("Now Stopping the Handler") handler.stop() print("Done!!") def __gen_entry(self): print("Generating entry") sz = len(self.__task['sources']) n = sz + 1 return { 'id': f"entry-{n}", 'rx': f"Rx{n}", 'fx': f"Fx{n}", }
关键总结
- 线程场景:直接修改实例属性即可,用锁保护共享访问避免竞态条件
- 进程场景:必须依赖IPC机制(队列、共享对象等),进程内存完全隔离,直接修改主进程实例属性对子进程无效
内容的提问来源于stack exchange,提问作者Oeg Bizz

