如何在ADF中记录日志到日志表并在阶段失败时发邮件?求Databricks通用PySpark代码
在Databricks中实现ADF执行阶段的日志记录与失败邮件通知
核心思路
可以通过封装通用PySpark工具函数,结合异常捕获机制,同时对接ADF传递的运行参数,实现执行细节日志记录与失败告警邮件推送。以下是可直接复用的代码实现:
1. 通用日志记录模块
首先需要在你的数据存储(如ADLS的Delta表、Azure SQL DB)中创建日志表,示例Delta表结构:
CREATE TABLE IF NOT EXISTS adf_execution_logs ( execution_id STRING, stage_name STRING, start_time TIMESTAMP, end_time TIMESTAMP, status STRING, error_message STRING, run_date DATE ) USING DELTA LOCATION '/path/to/adf-execution-logs';
然后在PySpark中封装日志写入函数:
from pyspark.sql import SparkSession import datetime def log_execution_stage(execution_id: str, stage_name: str, status: str, error_message: str = None): spark = SparkSession.getActiveSession() log_data = [ ( execution_id, stage_name, datetime.datetime.now(), datetime.datetime.now(), status, error_message, datetime.date.today() ) ] log_df = spark.createDataFrame( log_data, schema=["execution_id", "stage_name", "start_time", "end_time", "status", "error_message", "run_date"] ) # 追加模式写入日志表 log_df.write.mode("append").saveAsTable("adf_execution_logs")
2. 失败告警邮件模块
使用Python SMTP库对接邮件服务(如Azure SendGrid、企业SMTP服务器)实现告警推送:
import smtplib from email.mime.text import MIMEText from email.mime.multipart import MIMEMultipart import datetime def send_failure_alert(execution_id: str, stage_name: str, error_msg: str): # 替换为你的邮件服务配置 smtp_server = "smtp.sendgrid.net" smtp_port = 587 smtp_user = "apikey" smtp_password = "your-sendgrid-api-key" sender_email = "adf-alerts@yourdomain.com" receiver_emails = ["dev-ops@yourdomain.com"] # 构建邮件内容 subject = f"ADF执行失败告警:{stage_name} (Execution ID: {execution_id})" body = f""" ADF执行阶段失败详情: - 执行ID: {execution_id} - 失败阶段: {stage_name} - 错误信息: {error_msg} - 失败时间: {datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')} """ msg = MIMEMultipart() msg["From"] = sender_email msg["To"] = ", ".join(receiver_emails) msg["Subject"] = subject msg.attach(MIMEText(body, "plain")) # 发送邮件 try: with smtplib.SMTP(smtp_server, smtp_port) as server: server.starttls() server.login(smtp_user, smtp_password) server.send_message(msg) except Exception as e: print(f"告警邮件发送失败: {str(e)}")
3. 阶段执行的通用封装模板
在每个ADF调用的Databricks作业中,用以下模板包裹业务代码,自动完成日志记录与异常处理:
# 读取ADF传递的参数(通过Databricks作业参数传入) execution_id = dbutils.widgets.get("execution_id") stage_name = dbutils.widgets.get("stage_name") try: # 记录阶段启动日志 log_execution_stage(execution_id, stage_name, "RUNNING") # -------------------------- # 这里写入你的业务逻辑代码 # 示例:df = spark.read.table("source_table").write.mode("overwrite").saveAsTable("target_table") # -------------------------- # 记录阶段成功日志 log_execution_stage(execution_id, stage_name, "SUCCEEDED") except Exception as e: error_msg = str(e) # 记录阶段失败日志 log_execution_stage(execution_id, stage_name, "FAILED", error_msg) # 发送失败告警邮件 send_failure_alert(execution_id, stage_name, error_msg) # 抛出异常让ADF感知执行失败 raise e
与ADF集成的关键要点
- 在ADF的每个Databricks活动中,添加作业参数:
execution_id(使用ADF系统变量@pipeline().RunId)、stage_name(手动指定阶段名称,如"用户数据抽取阶段")。 - 确保Databricks集群有访问日志存储和SMTP服务器的网络权限。
- 若使用关系型数据库作为日志存储,需配置Databricks的JDBC连接权限。
内容的提问来源于stack exchange,提问作者Swati B
相关产品推荐
相关产品推荐

