Synapse Spark中Python日志写入ADLS文件的路径问题求助
解决Azure Synapse Spark中日志写入ADLS的路径问题
问题根源
Python标准库的logging.FileHandler仅支持本地文件系统路径,无法识别ADLS的abfss://协议。在Synapse Spark环境中,直接传入ADLS路径时,系统会自动将其视为本地路径并拼接默认挂载前缀(如/synfs/{job_id}/),导致路径无效,触发FileNotFoundError。
解决方案:自定义ADLS日志Handler
通过继承logging.Handler,使用Hadoop FileSystem API实现对ADLS路径的日志写入,绕过FileHandler的本地路径限制。
完整代码实现
import logging from pyspark.sql import SparkSession import sys class ADLSFileHandler(logging.Handler): def __init__(self, adls_path: str, mode: str = 'a'): super().__init__() self.adls_path = adls_path self.mode = mode # 获取当前活跃的Spark会话 self.spark = SparkSession.getActiveSession() if not self.spark: raise RuntimeError("未找到活跃的Spark会话") def emit(self, record): # 格式化日志消息 log_message = self.format(record) + '\n' # 通过Hadoop FileSystem操作ADLS hadoop_conf = self.spark._jsc.hadoopConfiguration() fs = self.spark._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf) path = self.spark._jvm.org.apache.hadoop.fs.Path(self.adls_path) if not fs.exists(path): # 文件不存在时创建并写入 output_stream = fs.create(path) output_stream.writeBytes(log_message.encode('utf-8')) output_stream.close() else: # 文件存在时追加写入 output_stream = fs.append(path) output_stream.writeBytes(log_message.encode('utf-8')) output_stream.close() def init_logger(name: str, logging_level: int = logging.DEBUG, adls_log_path: str = None) -> logging.Logger: _log_format = "%(levelname)s %(asctime)s %(name)s: %(message)s" _date_format = "%Y-%m-%d %I:%M:%S %p %z" _formatter = logging.Formatter(fmt=_log_format, datefmt=_date_format) _logger = logging.getLogger(name) _logger.setLevel(logging_level) _logger.propagate = False # 避免日志重复传播到根记录器 # 清空已有Handler,防止重复添加 if _logger.handlers: _logger.handlers = [] # 添加控制台输出Handler stream_handler = logging.StreamHandler(sys.stderr) stream_handler.setLevel(logging_level) stream_handler.setFormatter(_formatter) _logger.addHandler(stream_handler) # 添加ADLS日志Handler(如果传入路径) if adls_log_path: adls_handler = ADLSFileHandler(adls_log_path) adls_handler.setLevel(logging_level) adls_handler.setFormatter(_formatter) _logger.addHandler(adls_handler) return _logger # 使用示例 log_path = 'abfss:/container@storageaccountname.dfs.core.windows.net/Data/logging/data.log' logger = init_logger("error_logger", logging.ERROR, log_path) logger.error("测试ADLS错误日志写入")
关键说明
- 自定义Handler核心:重写
emit方法,利用Hadoop的FileSystem API直接操作ADLS路径,支持创建和追加两种写入模式。 - 避免重复日志:设置
_logger.propagate = False,防止日志被重复发送到根记录器的Handler。 - 路径兼容性:直接使用原生
abfss://路径,无需调整层级或添加前缀,系统会通过Spark的Hadoop配置识别ADLS协议。
替代方案:直接使用dbutils写入
如果不需要遵循logging模块的规范,也可以直接用dbutils.fs.put实现追加写入:
def write_log_to_adls(message: str, adls_path: str): # 用dbutils追加写入日志 dbutils.fs.put(adls_path, message + '\n', append=True) # 使用示例 write_log_to_adls("ERROR 2024-05-20 my_logger: 测试日志", log_path)
内容的提问来源于stack exchange,提问作者Morpheus273
相关产品推荐
相关产品推荐

