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

多进程环境下Queue日志记录失效问题排查

问题原因分析

你遇到的问题本质是Python多进程启动机制与logging模块的进程隔离特性导致的:

  1. 使用multiprocessing.Process类时,__init__方法在父进程中执行,你在这里配置的QueueHandler属于父进程的logger实例。
  2. 子进程启动后,logging系统是进程独立的:
    • 若用spawn模式(Windows、macOS 10.14+默认):子进程会重新导入模块,logging.Logger对象无法被正确序列化传递,子进程会重新创建同名但无handler的logger实例。
    • 若用fork模式(Unix/Linux默认):虽会复制父进程内存,但logger的状态在子进程中可能被重置,且跨进程共享handler存在潜在风险。
  3. __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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 14:41:01