使用Python multiprocessing.Pool时日志未输出到文件/控制台的问题
多进程+线程场景下Python日志不输出问题修复
问题原因分析
- 子进程重复初始化日志组件:在Windows(
spawn模式)或类Unix系统的进程启动逻辑中,子进程会重新导入logger_factory模块,导致LoggingManager的类变量被重置。子进程会创建独立的日志队列和监听进程,而非复用主进程的共享队列,最终子进程日志无法被主进程监听进程捕获。 - 跨进程传递不可序列化对象:原代码中向子进程传递
logger对象,但logging.Logger不支持序列化,会破坏日志流程(即使未显式报错)。 - 重复添加日志处理器:每次调用
get_logger都会新增QueueHandler,导致日志被多次发送到队列。
解决方案
- 主进程单例管控日志组件:仅在主进程初始化日志队列和监听进程,子进程直接使用主进程传递的共享队列。
- 日志器创建逻辑适配多进程:允许
get_logger接收外部队列,避免子进程触发新队列的初始化。 - 禁止跨进程传递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
相关产品推荐
相关产品推荐

