Azure Databricks任务:满足条件时发送邮件的实现方法
在Databricks中基于条件触发邮件通知
下面提供两种可行的实现方案,无需依赖Databricks Alerts,直接在笔记本代码中完成条件判断与邮件发送:
方案1:用Python SMTP库直接发送邮件
这是最直接的方式,只要Databricks集群能访问SMTP服务器即可实现。
代码示例
import smtplib from email.mime.text import MIMEText from email.mime.multipart import MIMEMultipart # 优化条件判断逻辑(避免重复计算) target_doc_count = df.filter(df.documenttype_id == 20).count() has_target_docs = target_doc_count > 0 if has_target_docs: # 配置SMTP参数(建议用Databricks Secrets存储敏感信息) smtp_server = "smtp.example.com" smtp_port = 587 sender_email = "your-alert@example.com" # 从Secrets获取密码,避免硬编码 sender_password = dbutils.secrets.get("your-secret-scope", "smtp-password") receiver_emails = ["user1@example.com", "user2@example.com"] # 构建邮件内容 msg = MIMEMultipart() msg['From'] = sender_email msg['To'] = ", ".join(receiver_emails) msg['Subject'] = "Databricks任务告警:检测到类型ID为20的文档" body = f""" 您好: 本次任务执行中,检测到存在documenttype_id为20的文档,共{target_doc_count}条。 请留意相关业务情况。 """ msg.attach(MIMEText(body, 'plain')) # 发送邮件 try: with smtplib.SMTP(smtp_server, smtp_port) as server: server.starttls() server.login(sender_email, sender_password) server.sendmail(sender_email, receiver_emails, msg.as_string()) print("告警邮件发送成功") except Exception as e: print(f"邮件发送失败:{str(e)}")
注意事项
- 确保集群安全组/网络规则允许访问SMTP服务器的对应端口
- 所有敏感信息(如SMTP密码)必须存储在Databricks Secrets中,禁止硬编码
方案2:调用Databricks Jobs API触发通知任务
如果希望解耦业务逻辑与邮件发送,可以创建一个专门负责发邮件的Databricks Job,在主笔记本条件满足时触发该任务。
步骤与代码示例
- 先创建一个独立的Databricks Job,任务内容为执行邮件发送代码(可复用方案1中的发送逻辑)
- 在主笔记本中触发该Job:
from databricks_cli.jobs.api import JobsApi from databricks_cli.sdk.api_client import ApiClient # 条件判断 target_doc_count = df.filter(df.documenttype_id == 20).count() has_target_docs = target_doc_count > 0 if has_target_docs: # 初始化Databricks API客户端 api_client = ApiClient( host=dbutils.notebook.entry_point.getDbutils().notebook().getContext().apiUrl().get(), token=dbutils.notebook.entry_point.getDbutils().notebook().getContext().apiToken().get() ) jobs_api = JobsApi(api_client) # 替换为你的通知Job ID alert_job_id = "1234567890" try: run_response = jobs_api.run_now(job_id=alert_job_id) print(f"已触发告警任务,运行ID:{run_response['run_id']}") except Exception as e: print(f"触发告警任务失败:{str(e)}")
注意事项
- 确保当前笔记本的执行权限允许调用Jobs API并触发目标Job
- 可在通知Job中通过传递参数的方式,动态传递告警详情(如文档数量)
最佳实践
- 避免重复计算:将条件判断中的统计结果存入变量,后续直接复用,减少Spark计算开销
- 异常捕获:在邮件发送或Job触发逻辑中添加异常处理,避免因通知失败导致整个任务中断
- 日志记录:关键步骤添加日志输出,方便后续排查问题
内容的提问来源于stack exchange,提问作者ImNotSureAboutStats
相关产品推荐
相关产品推荐

