使用多进程调度器的Dask结合Loguru时触发RuntimeError
Loguru + Dask Processes调度器日志报错解决
问题背景
启用Loguru的enqueue=True参数后,在Dask delayed函数内记录日志,使用processes调度器执行时会触发错误:
RuntimeError: SimpleQueue objects should only be shared between processes through inheritance
报错原因
Loguru开启enqueue后会创建SimpleQueue用于异步日志缓冲,但Dask的processes调度器会将任务函数及依赖(包括logger实例)序列化后传递给子进程。而SimpleQueue不支持跨进程序列化传递,仅能通过进程继承共享,因此触发该错误。
解决办法
方法1:子进程内单独初始化Loguru配置
不在全局提前设置带enqueue的处理器,让每个Dask任务在子进程内部自行添加日志处理器,确保每个子进程拥有独立的队列,避免跨进程共享问题:
import sys, dask from loguru import logger @dask.delayed def log(): # 子进程内部单独配置带enqueue的日志处理器 logger.add(sys.stderr, enqueue=True) logger.info("Logging!") dask.compute(*[log() for i in range(10)], scheduler="processes")
方法2:关闭Loguru的enqueue参数
若不需要异步日志的性能优化,直接移除enqueue=True配置,日志同步输出就不会创建SimpleQueue,自然避免序列化问题:
import sys, dask from loguru import logger logger.add(sys.stderr) # 移除enqueue=True配置 @dask.delayed def log(): logger.info("Logging!") dask.compute(*[log() for i in range(10)], scheduler="processes")
方法3:对接Dask原生日志系统
将Loguru的日志转发到Dask原生日志流,适配分布式任务的日志处理逻辑:
import dask import logging from loguru import logger # 自定义处理器,将Loguru日志转发给Dask的root logger class DaskLogHandler(logging.Handler): def emit(self, record): logging.getLogger().handle(record) logger.add(DaskLogHandler(), format="{message}") @dask.delayed def log(): logger.info("Logging!") dask.compute(*[log() for i in range(10)], scheduler="processes")
内容的提问来源于stack exchange,提问作者David Hoffman
相关产品推荐
相关产品推荐

