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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.01 20:27:26