如何在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()
关键说明
- 环境变量传递文件名是因为Dask Worker进程会继承主进程的环境变量,无需手动传递参数。
- 配置Logger时先清空现有Handler,避免Worker重启或重复配置时出现重复日志输出。
- 日志格式加入
%(process)d可以区分不同Worker进程的日志,方便排查问题。
内容的提问来源于stack exchange,提问作者Lei Yu
相关产品推荐
相关产品推荐

