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

类实例中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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 19:44:55