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

使用Python multiprocessing.Pool时日志未输出到文件/控制台的问题

多进程+线程场景下Python日志不输出问题修复

问题原因分析

  1. 子进程重复初始化日志组件:在Windows(spawn模式)或类Unix系统的进程启动逻辑中,子进程会重新导入logger_factory模块,导致LoggingManager的类变量被重置。子进程会创建独立的日志队列和监听进程,而非复用主进程的共享队列,最终子进程日志无法被主进程监听进程捕获。
  2. 跨进程传递不可序列化对象:原代码中向子进程传递logger对象,但logging.Logger不支持序列化,会破坏日志流程(即使未显式报错)。
  3. 重复添加日志处理器:每次调用get_logger都会新增QueueHandler,导致日志被多次发送到队列。

解决方案

  1. 主进程单例管控日志组件:仅在主进程初始化日志队列和监听进程,子进程直接使用主进程传递的共享队列。
  2. 日志器创建逻辑适配多进程:允许get_logger接收外部队列,避免子进程触发新队列的初始化。
  3. 禁止跨进程传递Logger对象:子进程内部通过名称+传入队列初始化日志器,规避序列化问题。

修正后的代码

logger_factory.py

import logging
from logging import handlers
import multiprocessing as mp
import datetime
import os

class LoggingManager:
    _log_queue = None
    _listener_process = None
    _is_main_process = True

    @classmethod
    def initialize(cls):
        """仅主进程调用,初始化日志队列与监听进程"""
        if cls._log_queue is None:
            cls._log_queue = mp.Queue()
            cls._start_listener_process()
        cls._is_main_process = True

    @classmethod
    def get_log_queue(cls):
        """获取主进程初始化的日志队列"""
        return cls._log_queue

    @classmethod
    def _start_listener_process(cls):
        """启动监听进程(仅主进程执行)"""
        if cls._listener_process is None:
            cls._listener_process = mp.Process(target=cls._listener_process_target, args=(cls._log_queue,))
            cls._listener_process.start()

    @classmethod
    def _listener_process_target(cls, log_queue):
        """监听进程消费日志逻辑"""
        print("Listener process started.")
        cls.listener_configurer()

        while True:
            try:
                record = log_queue.get()
                if record == "STOP":
                    print("Stopping listener process.")
                    break

                logger = logging.getLogger(record.name)
                logger.handle(record)

            except Exception as e:
                print(f"Error in listener process: {e}")
                break

    @classmethod
    def listener_configurer(cls):
        """配置日志输出处理器"""
        log_file_path = 'C:/Logs/test_{0}.log'.format(datetime.datetime.now().strftime('%Y%m%d%H%M%S'))
        os.makedirs(os.path.dirname(log_file_path), exist_ok=True)

        root = logging.getLogger()
        root.handlers.clear()  # 清空默认处理器,避免重复输出

        # 文件处理器:记录INFO及以上级别
        file_handler = handlers.TimedRotatingFileHandler(
            log_file_path, when="midnight", interval=1, backupCount=5
        )
        file_handler.setLevel(logging.INFO)

        # 控制台处理器:记录所有级别
        console_handler = logging.StreamHandler()
        console_handler.setLevel(logging.DEBUG)

        formatter = logging.Formatter('%(asctime)s %(processName)-10s %(name)s %(levelname)-8s %(message)s')
        file_handler.setFormatter(formatter)
        console_handler.setFormatter(formatter)

        root.addHandler(file_handler)
        root.addHandler(console_handler)
        root.setLevel(logging.DEBUG)

    @classmethod
    def get_logger(cls, name, log_queue=None):
        """获取日志器,支持传入外部队列(子进程使用)"""
        logger = logging.getLogger(name)
        logger.setLevel(logging.DEBUG)

        # 避免重复添加QueueHandler
        has_queue_handler = any(isinstance(handler, handlers.QueueHandler) for handler in logger.handlers)
        if not has_queue_handler:
            target_queue = log_queue or cls._log_queue
            if target_queue:
                queue_handler = handlers.QueueHandler(target_queue)
                logger.addHandler(queue_handler)

        return logger

    @classmethod
    def stop_listener(cls):
        """优雅停止监听进程"""
        if cls._log_queue:
            cls._log_queue.put("STOP")
        if cls._listener_process:
            cls._listener_process.join()

    @classmethod
    def set_as_child_process(cls):
        """标记当前进程为子进程,禁止初始化新日志组件"""
        cls._is_main_process = False

automation.py

import concurrent.futures
import multiprocessing as mp
import sys
import time
from logger_factory import LoggingManager

def thread_process_client_query_exec(log_q):
    logger = LoggingManager.get_logger("thread_process_client_query_exec", log_q)
    logger.info("Started thread for query execution.")
    print(f"Queue size in thread after logging: {log_q.qsize()}")
    time.sleep(1)
    logger.info("Completed thread for query execution.")

def process_fun(log_q):
    # 标记为子进程,避免创建新日志组件
    LoggingManager.set_as_child_process()
    logger = LoggingManager.get_logger("process_fun", log_q)
    logger.info("Started process_fun.")
    print(f"Queue size in process_fun after logging: {log_q.qsize()}")
    sys.stdout.flush()
    future_to_query = []
    with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor:
        for query in range(5):
            res = executor.submit(thread_process_client_query_exec, log_q)
            future_to_query.append(res)

        for future in concurrent.futures.as_completed(future_to_query):
            try:
                future.result()
            except Exception as exc:
                logger.error(f"Exception in thread: {exc}")

def do_run_process_automation_test(log_q):
    logger = LoggingManager.get_logger(__name__, log_q)
    logger.info("Started do_run_automation_test.")
    
    with mp.Pool(2) as pool:
        results = []
        for i in range(5):
            # 仅传递日志队列,不传递logger对象
            p = pool.apply_async(process_fun, (log_q,))
            results.append(p)

        pool.close()
        pool.join()
    logger.info("Completed do_run_automation_test.")

main.py

import time
from logger_factory import LoggingManager
import automation

def main():
    # 主进程初始化日志系统
    LoggingManager.initialize()
    log_q = LoggingManager.get_log_queue()

    logger = LoggingManager.get_logger(__name__)
    logger.info("Logging from the main process")
    logger.debug("Debug message from main process")
    logger.error("Error message from main process")

    automation.do_run_process_automation_test(log_q)

    print(f"Queue size after all tasks: {log_q.qsize()}")

    time.sleep(2)

    LoggingManager.stop_listener()

if __name__ == "__main__":
    main()

关键修改说明

  • 新增initialize方法,确保日志队列和监听进程仅在主进程初始化一次。
  • get_logger支持传入外部队列,子进程直接使用主进程传递的队列发送日志。
  • 新增set_as_child_process标记,避免子进程误创建独立的日志组件。
  • 移除跨进程传递logger对象的逻辑,规避序列化问题。
  • 在日志配置中清空根日志器现有处理器,避免重复输出日志。

内容的提问来源于stack exchange,提问作者Basavaraj Biradar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 05:05:00