多进程环境下Queue日志记录失效问题排查
问题原因分析
你遇到的问题本质是Python多进程启动机制与logging模块的进程隔离特性导致的:
- 使用
multiprocessing.Process类时,__init__方法在父进程中执行,你在这里配置的QueueHandler属于父进程的logger实例。 - 子进程启动后,logging系统是进程独立的:
- 若用
spawn模式(Windows、macOS 10.14+默认):子进程会重新导入模块,logging.Logger对象无法被正确序列化传递,子进程会重新创建同名但无handler的logger实例。 - 若用
fork模式(Unix/Linux默认):虽会复制父进程内存,但logger的状态在子进程中可能被重置,且跨进程共享handler存在潜在风险。
- 若用
__init__中的logger.info能写入日志,是因为这条日志在父进程执行,通过父进程的QueueHandler发送到队列;而run方法在子进程执行,此时子进程的logger没有绑定handler,自然无法输出日志。
基于Process类的优化方案
核心思路:把logger的配置逻辑移到子进程执行的run方法中,确保handler在子进程内创建并绑定到子进程的logger实例。
修改后的代码示例:
import logging import random import sys import time from logging import handlers from multiprocessing import Queue, Process from queue import Empty, Full def worker_configurer(log_queue, idx): logger = logging.getLogger(".".join(("A", "worker", str(idx)))) # 清空可能存在的handler,避免重复添加 logger.handlers.clear() h = handlers.QueueHandler(log_queue) logger.addHandler(h) logger.setLevel(logging.INFO) print( f"configured worker {idx} with logger {logger.name} with handlers: {logger.handlers.copy()}" ) return logger class Worker(Process): worker_idx = 0 def __init__(self, work_queue, log_queue, **kwargs): super(Worker, self).__init__() self.idx = Worker.worker_idx Worker.worker_idx += 1 # 仅保存必要参数,不在父进程配置logger self.log_queue = log_queue self.work_queue = work_queue print(f"Worker {self.idx} initialized in parent process") def run(self): # 在子进程内配置logger self.logger = worker_configurer(self.log_queue, self.idx) print( f"(inside run): self.logger name: {self.logger.name}, handlers:" f" {self.logger.handlers.copy()}" ) self.logger.info(f"worker {self.idx} started!") # 现在会正常输出到日志 while True: job_duration = self.work_queue.get() if job_duration is None: print(f"Worker {self.idx} received stop signal") break time.sleep(job_duration) print(f"worker {self.idx} finished job of length {job_duration}") self.logger.info(f"worker {self.idx} finished job of length {job_duration}") def listener_configurer(): logging.basicConfig( filename="mp_log.log", filemode="a", format="%(asctime)s | %(name)s | %(levelname)s | %(message)s", datefmt="%d-%b-%y %H:%M:%S", level=logging.INFO, ) def listener_process(queue, configurer): configurer() logger = logging.getLogger("A") while True: try: record = queue.get(timeout=5) print("bagged a message from the queue") if record is None: break logger = logging.getLogger(record.name) logger.handle(record) except Empty: pass except Exception: import traceback print("Whoops! Problem:", file=sys.stderr) traceback.print_exc(file=sys.stderr) if __name__ == "__main__": listener_configurer() logger = logging.getLogger("A") logger.warning("Logger Active!") work_queue = Queue(5) log_queue = Queue(100) listener = Process(target=listener_process, args=(log_queue, listener_configurer)) listener.start() num_workers = 2 workers = [] for i in range(num_workers): w = Worker( work_queue, log_queue=log_queue, ) w.start() workers.append(w) logger.info(f"worker {i} created") num_jobs = 10 jobs_assigned = 0 while jobs_assigned < num_jobs: try: work_queue.put(random.random() * 2, timeout=0.1) jobs_assigned += 1 except Full: pass print("Call it a day and send stop sentinel to everybody") for i in range(num_workers): work_queue.put(None) log_queue.put(None) for w in workers: w.join() print("another worker retired!") listener.join()
验证结果
修改后控制台会输出:
Worker 0 initialized in parent process Worker 1 initialized in parent process bagged a message from the queue bagged a message from the queue configured worker 1 with logger A.worker.1 with handlers: [<QueueHandler (NOTSET)>] (inside run): self.logger name: A.worker.1, handlers: [<QueueHandler (NOTSET)>] configured worker 0 with logger A.worker.0 with handlers: [<QueueHandler (NOTSET)>] (inside run): self.logger name: A.worker.0, handlers: [<QueueHandler (NOTSET)>] bagged a message from the queue bagged a message from the queue worker 0 finished job of length 0.xxxxxx ...
日志文件中会包含worker x started!以及所有任务完成的日志条目,符合预期。
额外注意事项
- 禁止在父进程中给子进程的logger添加handler,所有子进程的日志配置必须在子进程内部完成。
- 即使使用
fork模式,仍建议在子进程中重新配置logger,避免跨进程的handler冲突风险。
内容的提问来源于stack exchange,提问作者AirSquid
相关产品推荐
相关产品推荐

