类实例中ProcessPoolExecutor任务提交报错求助:无法序列化线程锁对象
多进程与线程交互中的任务提交问题
需求说明
需要实现一个类,具备以下功能:
- 内部通过线程每秒更新
self.current_time; - 每次更新时间时,自动将处理任务(示例中为
print_time,实际为复杂计算)提交给ProcessPoolExecutor多进程执行,而非由主程序定时触发。
初始代码(存在触发时机问题)
这段代码能运行,但任务执行时机由主程序控制,不符合需求:
import time import threading import concurrent.futures class time_print_class: def __init__(self, tid): self.tid = tid time_thread = threading.Thread(target=self.update_time) time_thread.start() self.current_time = 0 def update_time(self): while True: self.current_time = time.time() time.sleep(1) def print_time(self): # 实际为复杂任务,需多进程执行 print(self.current_time, "on", self.tid) if __name__ == "__main__": # 初始化60个实例 time_print_insts = [time_print_class(n) for n in range(60)] # 主程序定时触发任务 executor = concurrent.futures.ProcessPoolExecutor() while True: for tpi in time_print_insts: future = executor.submit(tpi.print_time) time.sleep(1)
尝试的错误代码及问题
尝试让类内部提交任务,但触发cannot pickle '_thread.lock' object错误:
import time import threading import concurrent.futures class time_print_class: def __init__(self, tid, executor): self.tid = tid self.executor = executor time_thread = threading.Thread(target=self.update_time) time_thread.start() self.current_time = 0 def update_time(self): while True: self.current_time = time.time() self.future = executor.submit(self.print_time) time.sleep(1) def print_time(self): print(self.current_time, "on", self.tid) if __name__ == "__main__": executor = concurrent.futures.ProcessPoolExecutor() time_print_insts = [time_print_class(n, executor) for n in range(60)]
错误原因:ProcessPoolExecutor提交任务时需要序列化(pickle)待执行对象,而类实例包含线程锁、executor等无法被pickle的对象,直接提交实例方法会导致整个实例被序列化,触发报错。
解决方案
方案1:将任务改为无状态函数,传递必要参数
把print_time改成独立函数,提交任务时只传递需要的current_time和tid,避免序列化整个类实例:
import time import threading import concurrent.futures def process_time(current_time, tid): # 实际复杂逻辑写在这里 print(current_time, "on", tid) class time_print_class: def __init__(self, tid, executor): self.tid = tid self.executor = executor time_thread = threading.Thread(target=self.update_time) time_thread.start() self.current_time = 0 def update_time(self): while True: self.current_time = time.time() # 提交独立函数,只传必要数据 self.executor.submit(process_time, self.current_time, self.tid) time.sleep(1) if __name__ == "__main__": executor = concurrent.futures.ProcessPoolExecutor() time_print_insts = [time_print_class(n, executor) for n in range(60)] # 保持主进程运行 try: while True: time.sleep(3600) except KeyboardInterrupt: executor.shutdown()
方案2:使用队列做任务中转(解耦类与Executor)
类的update_time只负责将任务数据放入队列,主进程单独开线程从队列取数据并提交给Executor,彻底避免类持有Executor带来的序列化问题:
import time import threading import concurrent.futures from queue import Queue def process_time(current_time, tid): # 实际复杂逻辑写在这里 print(current_time, "on", tid) def queue_worker(queue, executor): while True: task_data = queue.get() if task_data is None: break executor.submit(process_time, *task_data) queue.task_done() class time_print_class: def __init__(self, tid, task_queue): self.tid = tid self.task_queue = task_queue time_thread = threading.Thread(target=self.update_time) time_thread.start() self.current_time = 0 def update_time(self): while True: self.current_time = time.time() # 将任务数据放入队列 self.task_queue.put((self.current_time, self.tid)) time.sleep(1) if __name__ == "__main__": task_queue = Queue() executor = concurrent.futures.ProcessPoolExecutor() # 启动队列消费线程 threading.Thread(target=queue_worker, args=(task_queue, executor), daemon=True).start() # 初始化实例 time_print_insts = [time_print_class(n, task_queue) for n in range(60)] # 保持主进程运行 try: while True: time.sleep(3600) except KeyboardInterrupt: # 停止队列 worker task_queue.put(None) executor.shutdown()
方案优势对比:
- 方案1实现简单,适合任务逻辑与类耦合度低的场景;
- 方案2解耦了类与Executor,更灵活,适合需要控制任务流量、添加重试机制等复杂场景。
内容的提问来源于stack exchange,提问作者wanga10000
相关产品推荐
相关产品推荐

