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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 23:09:50