如何使用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
相关产品推荐
相关产品推荐

