Databricks Python笔记本日志配置:对接Azure Log Analytics
在Databricks Python笔记本中实现全局日志发送至Azure Log Analytics
1. 基于Python标准库封装自定义日志处理器
利用Python内置的logging模块扩展,实现适配Azure Log Analytics的日志处理器,自动注入笔记本名称、作业启动时间等元数据:
import logging import requests import json from datetime import datetime import dbutils import hmac import hashlib import base64 class AzureLogAnalyticsHandler(logging.Handler): def __init__(self, workspace_id, shared_key, log_type): super().__init__() self.workspace_id = workspace_id self.shared_key = shared_key self.log_type = log_type # 预加载Databricks会话元数据 ctx = dbutils.notebook.entry_point.getDbutils().notebook().getContext() self.notebook_name = ctx.notebookPath().get() self.job_start_time = datetime.now().isoformat() self.cluster_id = ctx.clusterId().get() self.job_id = ctx.tags().get("jobId").orElse("interactive_session") def _generate_signature(self, date_str, content_len): # 生成Azure Log Analytics API要求的身份签名 string_to_sign = f"POST\n{content_len}\napplication/json\nx-ms-date:{date_str}" decoded_key = base64.b64decode(self.shared_key) signature = hmac.new(decoded_key, string_to_sign.encode('utf-8'), hashlib.sha256).digest() return f"SharedKey {self.workspace_id}:{base64.b64encode(signature).decode()}" def emit(self, record): # 构造包含元数据的日志条目 log_item = { "message": self.format(record), "log_level": record.levelname, "notebook_path": self.notebook_name, "job_start_time": self.job_start_time, "cluster_id": self.cluster_id, "job_id": self.job_id, "utc_timestamp": datetime.utcnow().isoformat() + "Z" } # 发送日志到Azure Log Analytics date_str = datetime.utcnow().strftime('%a, %d %b %Y %H:%M:%S GMT') payload = json.dumps([log_item]) signature = self._generate_signature(date_str, len(payload)) headers = { "Content-Type": "application/json", "Authorization": signature, "Log-Type": self.log_type, "x-ms-date": date_str } api_url = f"https://{self.workspace_id}.ods.opinsights.azure.com/api/logs?api-version=2016-04-01" try: response = requests.post(api_url, headers=headers, data=payload) response.raise_for_status() except Exception: # 发送失败时回退到本地日志记录 self.handleError(record)
2. 初始化全局日志实例
在笔记本开头执行以下代码,创建全局可用的日志对象,后续直接调用标准日志方法:
# 替换为你的Azure Log Analytics配置 WORKSPACE_ID = "你的工作区ID" # 建议用Databricks Secrets管理密钥,避免硬编码 SHARED_KEY = dbutils.secrets.get("你的密钥范围", "log_analytics_shared_key") LOG_TYPE = "DatabricksNotebookLogs" # 全局日志实例 log = logging.getLogger("databricks_azure_logger") log.setLevel(logging.DEBUG) # 避免重复注册处理器 if not log.handlers: formatter = logging.Formatter('%(asctime)s - %(module)s - %(message)s') az_handler = AzureLogAnalyticsHandler(WORKSPACE_ID, SHARED_KEY, LOG_TYPE) az_handler.setFormatter(formatter) log.addHandler(az_handler)
3. 日志使用示例
直接调用熟悉的日志方法,日志会自动携带元数据发送到Azure Log Analytics:
log.info("笔记本作业开始执行") log.warn("please take care...") log.error("数据处理失败", exc_info=True)
关键注意事项
- 元数据扩展:可根据需求添加更多元数据,如当前用户
ctx.userName().get()、集群名称ctx.tags().get("clusterName").orElse(None)等。 - 性能优化:高频率日志场景下,可修改处理器实现批量发送逻辑,或结合
logging.QueueHandler实现异步发送,避免阻塞业务代码。 - 权限控制:确保Databricks集群有访问Azure Log Analytics API的网络权限,若使用VNet隔离,需配置相应的网络规则。
内容的提问来源于stack exchange,提问作者ndueck
相关产品推荐
相关产品推荐

