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

如何在Dask中有效配置Python Logger?多进程日志方案失效

Dask多进程日志配置失效的解决方法

问题根源

Dask的Worker进程是独立启动的,主进程中配置好的Logger对象无法被序列化传递给Worker,导致Worker进程的Logger没有配置,自然不会输出日志。

方案一:通过环境变量传递日志文件名,在Worker任务中自动配置

这是对临时方案的优化,无需手动传递文件名,通过环境变量让Worker自动获取:

import logging
import os
from datetime import datetime
import dask
from dask.distributed import LocalCluster, Client

def setup_logger(log_filename=None):
    # 优先从环境变量获取文件名
    if log_filename is None:
        log_filename = os.environ.get('DASK_LOG_FILENAME')
    logger = logging.getLogger()
    # 避免重复添加Handler(Worker进程重启时可能重复执行)
    for handler in list(logger.handlers):
        logger.removeHandler(handler)
    logger.setLevel(logging.INFO)
    # 配置文件输出
    file_handler = logging.FileHandler(log_filename)
    # 配置控制台输出
    stream_handler = logging.StreamHandler()
    # 统一日志格式,加上进程ID区分不同Worker
    formatter = logging.Formatter('%(asctime)s - %(process)d - %(levelname)s - %(message)s')
    file_handler.setFormatter(formatter)
    stream_handler.setFormatter(formatter)
    logger.addHandler(file_handler)
    logger.addHandler(stream_handler)
    return logger

def fake_task():
    # 每个任务直接初始化Logger(自动从环境变量拿文件名)
    log = setup_logger()
    log.info('任务执行中...')

if __name__ == '__main__':
    # 生成带实时信息的日志文件名
    log_filename = datetime.now().strftime('%I%M%p.log').lower()  # 例如0830am.log
    # 设置环境变量,Worker进程会继承这个变量
    os.environ['DASK_LOG_FILENAME'] = log_filename
    # 主进程自身的Logger配置
    main_logger = setup_logger(log_filename)
    main_logger.info('主进程启动')

    # 启动Dask集群
    cluster = LocalCluster(n_workers=10, threads_per_worker=1)
    client = Client(cluster)

    # 生成任务
    tasks = [dask.delayed(fake_task)() for _ in range(10)]
    # 执行任务
    dask.delayed(len)(tasks).compute()

方案二:利用Dask Worker启动钩子全局配置Logger

这种方式更高效,仅在每个Worker启动时配置一次Logger,无需每个任务都初始化:

import logging
import os
from datetime import datetime
import dask
from dask.distributed import LocalCluster, Client

def setup_logger(log_filename=None):
    logger = logging.getLogger()
    for handler in list(logger.handlers):
        logger.removeHandler(handler)
    logger.setLevel(logging.INFO)
    file_handler = logging.FileHandler(log_filename)
    stream_handler = logging.StreamHandler()
    formatter = logging.Formatter('%(asctime)s - %(process)d - %(levelname)s - %(message)s')
    file_handler.setFormatter(formatter)
    stream_handler.setFormatter(formatter)
    logger.addHandler(file_handler)
    logger.addHandler(stream_handler)
    return logger

def setup_worker_logger():
    # Worker启动时执行,从环境变量获取日志文件名
    log_filename = os.environ.get('DASK_LOG_FILENAME')
    setup_logger(log_filename)

def fake_task():
    # 直接使用全局配置好的Logger
    log = logging.getLogger()
    log.info('任务执行中...')

if __name__ == '__main__':
    log_filename = datetime.now().strftime('%I%M%p.log').lower()
    os.environ['DASK_LOG_FILENAME'] = log_filename
    main_logger = setup_logger(log_filename)
    main_logger.info('主进程启动')

    # 启动集群时,注册Worker启动钩子
    cluster = LocalCluster(n_workers=10, threads_per_worker=1)
    client = Client(cluster)
    # 注册Worker启动时要执行的函数
    client.register_worker_callbacks(setup=setup_worker_logger)

    tasks = [dask.delayed(fake_task)() for _ in range(10)]
    dask.delayed(len)(tasks).compute()

关键说明

  1. 环境变量传递文件名是因为Dask Worker进程会继承主进程的环境变量,无需手动传递参数。
  2. 配置Logger时先清空现有Handler,避免Worker重启或重复配置时出现重复日志输出。
  3. 日志格式加入%(process)d可以区分不同Worker进程的日志,方便排查问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 16:25:13