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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 07:45:52