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

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,在主笔记本条件满足时触发该任务。

步骤与代码示例

  1. 先创建一个独立的Databricks Job,任务内容为执行邮件发送代码(可复用方案1中的发送逻辑)
  2. 在主笔记本中触发该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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 13:15:25