使用multiprocessing.Pool结合类封装Logger时遇无法pickle _thread.lock对象问题
多进程类实例序列化与日志处理问题
问题描述
尝试将数据处理任务并行化,把逻辑封装在DataProcessor类中,每个实例需要向集中日志文件记录进度,但调用multiprocessing.Pool.map()时出现序列化错误。
原代码如下:
import logging import multiprocessing class DataProcessor: def __init__(self, name): self.name = name self.logger = logging.getLogger("ProcessorLogger") self.logger.setLevel(logging.INFO) def run(self, data): self.logger.info(f"{self.name} is processing {data}") return data * 2 def worker(obj_and_data): obj, data = obj_and_data return obj.run(data) if __name__ == "__main__": tasks = [(DataProcessor(f"Proc-{i}"), i) for i in range(5)] with multiprocessing.Pool(processes=4) as pool: results = pool.map(worker, tasks) print(results)
运行后报错:
Traceback (most recent call last): File "script.py", line 22, in <module> results = pool.map(worker, tasks) ... File "/usr/lib/python3.10/multiprocessing/reduction.py", line 51, in dump ForkingPickler(file, protocol).dump(obj) TypeError: cannot pickle '_thread.lock' object AttributeError: Can't pickle _thread.lock objects
已尝试方案:
- 将logger设为全局变量,但不同处理器需要不同配置,不适用
- 使用
pathos.multiprocessing(基于dill),仍存在序列化问题或日志无法正常输出
错误原因
logging模块的Logger对象内部包含线程锁(_thread.lock),用于保证多线程环境下日志输出的安全性。而multiprocessing.Pool在传递任务参数时,需要对对象进行序列化(pickle),但pickle无法序列化这类底层的锁对象,因此当类实例持有Logger时,序列化过程会失败。
正确的日志处理方式
1. 在子进程中初始化Logger
不在主进程的类构造方法中创建Logger,而是在子进程执行的run方法内初始化Logger。这样每个子进程会独立创建自己的Logger实例,避免了序列化锁对象的问题,同时还能根据实例的属性配置不同的日志标识。
修改后的代码示例:
import logging import multiprocessing class DataProcessor: def __init__(self, name): self.name = name # 不在主进程初始化Logger self.logger = None def run(self, data): # 子进程执行时初始化Logger if not self.logger: self.logger = logging.getLogger(f"ProcessorLogger.{self.name}") self.logger.setLevel(logging.INFO) # 添加文件Handler,确保所有进程日志写入同一文件 handler = logging.FileHandler("processing.log") formatter = logging.Formatter("%(asctime)s - %(name)s - %(levelname)s - %(message)s") handler.setFormatter(formatter) # 避免重复添加Handler(子进程可能多次调用run时) if not self.logger.handlers: self.logger.addHandler(handler) self.logger.info(f"{self.name} is processing {data}") return data * 2 def worker(obj_and_data): obj, data = obj_and_data return obj.run(data) if __name__ == "__main__": tasks = [(DataProcessor(f"Proc-{i}"), i) for i in range(5)] with multiprocessing.Pool(processes=4) as pool: results = pool.map(worker, tasks) print(results)
2. 全局配置日志,子进程自动继承
在主进程中完成日志的全局配置(比如设置Handler、Formatter、级别等),子进程启动时会自动继承这些配置,此时不需要在类实例中持有Logger,直接在run方法中获取全局Logger即可:
import logging import multiprocessing # 主进程中全局配置日志 logging.basicConfig( level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s", handlers=[logging.FileHandler("processing.log")] ) class DataProcessor: def __init__(self, name): self.name = name def run(self, data): # 直接获取全局配置的Logger logger = logging.getLogger(f"ProcessorLogger.{self.name}") logger.info(f"{self.name} is processing {data}") return data * 2 def worker(obj_and_data): obj, data = obj_and_data return obj.run(data) if __name__ == "__main__": tasks = [(DataProcessor(f"Proc-{i}"), i) for i in range(5)] with multiprocessing.Pool(processes=4) as pool: results = pool.map(worker, tasks) print(results)
关键注意点
- 避免在主进程创建的对象中持有无法序列化的资源(如锁、文件句柄、网络连接等),这些资源都不能被pickle序列化传递给子进程。
- 多进程日志写入同一文件时,
logging模块的FileHandler默认会处理多进程下的文件锁(取决于操作系统),如果需要更严格的控制,可以使用QueueHandler将日志消息发送到主进程统一处理。
内容的提问来源于stack exchange,提问作者ssd
相关产品推荐
相关产品推荐

