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

如何使用multiprocessing.Pool.apply_async实现单文件日志记录

解决multiprocessing.Pool下的单文件日志输出问题

你遇到的核心问题有两个:一是multiprocessing.Pool的子进程是预创建并复用的,直接通过apply_async传递队列和配置函数无法保证每个子进程都正确初始化日志;二是你在任务还没全部执行完就给队列发了None,导致监听进程提前退出,丢失后续日志。

下面是修正后的可运行版本,同时解决主进程日志捕获的问题:

import logging
import logging.handlers
import numpy as np
import time
import multiprocessing
import pandas as pd

log_file = 'log_file.log'

# 监听进程配置
def listener_configurer():
    root = logging.getLogger()
    root.setLevel(logging.DEBUG)  # 确保能接收所有级别日志
    h = logging.FileHandler(log_file)
    f = logging.Formatter('%(asctime)s %(processName)-10s %(name)s %(levelname)-8s %(message)s')
    h.setFormatter(f)
    root.addHandler(h)

def listener_process(queue):
    listener_configurer()
    while True:
        try:
            record = queue.get()
            if record is None:
                break
            logger = logging.getLogger(record.name)
            logger.handle(record)
        except Exception:
            import sys, traceback
            print('Whoops! Problem:', file=sys.stderr)
            traceback.print_exc(file=sys.stderr)

# 子进程日志初始化:给Pool的每个子进程统一配置
def worker_init(queue):
    h = logging.handlers.QueueHandler(queue)
    root = logging.getLogger()
    root.addHandler(h)
    root.setLevel(logging.DEBUG)
    # 移除默认的StreamHandler,避免重复输出到控制台
    for handler in list(root.handlers):
        if isinstance(handler, logging.StreamHandler):
            root.removeHandler(handler)

# 工作函数:不需要再传queue和configurer,因为子进程已经初始化好了
def worker_function(sleep_time, name):
    start_message = 'Worker {} started and will now sleep for {}s'.format(name, sleep_time)
    logging.info(start_message)
    time.sleep(sleep_time)
    success_message = 'Worker {} has finished sleeping for {}s'.format(name, sleep_time)
    logging.info(success_message)

def main_with_pool():
    start_time = time.time()
    # 使用Manager.Queue避免跨进程队列的潜在问题(Pool子进程和主进程通信更可靠)
    queue = multiprocessing.Manager().Queue(-1)
    listener = multiprocessing.Process(target=listener_process, args=(queue,))
    listener.start()

    # 初始化Pool时指定子进程的初始化函数和参数
    with multiprocessing.Pool(processes=3, initializer=worker_init, initargs=(queue,)) as pool:
        job_list = [np.random.randint(10) / 2 for i in range(10)]
        single_thread_time = np.sum(job_list)
        # 提交所有任务
        results = []
        for i, sleep_time in enumerate(job_list):
            name = str(i)
            results.append(pool.apply_async(worker_function, args=(sleep_time, name)))
        
        # 等待所有任务完成
        for res in results:
            res.get()
    
    # 所有任务完成后,发送终止信号给监听进程
    queue.put_nowait(None)
    listener.join()

    # 主进程日志也写入文件:配置主进程的QueueHandler
    root = logging.getLogger()
    root.addHandler(logging.handlers.QueueHandler(queue))
    root.setLevel(logging.INFO)
    final_message = "Script execution time was {:.2f}s, but single-thread time was {:.2f}s".format(
        (time.time() - start_time), single_thread_time )
    logging.info(final_message)
    print(final_message)

if __name__ == "__main__":
    main_with_pool()

关键改动说明:

  • Pool子进程统一初始化:用Pool(initializer=worker_init, initargs=(queue,))让每个预创建的子进程启动时就配置好日志队列,避免每次任务传递参数,同时适配Pool的进程复用机制。
  • 使用Manager.Queue:相比普通的multiprocessing.Queue,Manager创建的队列在Pool场景下跨进程通信更稳定,避免潜在的死锁或数据丢失。
  • 等待所有任务完成再终止监听:通过res.get()等待每个apply_async任务完成,再向队列发送None,确保所有日志都被监听进程处理。
  • 主进程日志捕获:主进程也配置QueueHandler,这样主进程的日志也会被写入同一个文件。

要不要改用Celery?

这取决于你的需求:

  • 如果只是本地并发任务管理,不需要分布式、任务持久化或复杂的调度,上面的Pool方案完全足够,轻量且不需要额外依赖。
  • 如果你的任务量极大(比如数千个以上)、需要跨机器分布式执行、任务失败重试、定时任务等高级特性,那么Celery是更好的选择,但需要额外配置消息中间件(如Redis、RabbitMQ),学习成本更高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:53:09