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

为Google Dataproc上的PySpark作业配置邮件告警代码实现方案

在Google Dataproc的PySpark作业中实现邮件告警

1. 敏感信息与SMTP配置管理

直接在代码中硬编码SMTP账号密码存在安全风险,推荐用Google Cloud Secret Manager存储敏感信息:

  • 创建Secret(例如命名为smtp-credentials),内容为JSON格式:
    {"smtp_host": "smtp.gmail.com", "smtp_port": 587, "smtp_user": "your-alert-email@gmail.com", "smtp_pass": "your-app-specific-password"}
    
  • 给Dataproc集群的服务账号添加roles/secretmanager.secretAccessor权限,确保作业能读取该Secret。
  • 若集群未预装google-cloud-secret-manager库,初始化集群时通过以下命令安装:
    gcloud dataproc clusters create your-cluster-name \
      --initialization-actions gs://dataproc-initialization-actions/python/pip-install.sh \
      --metadata PIP_PACKAGES=google-cloud-secret-manager
    

2. 编写邮件发送工具函数

在PySpark代码中实现通用的告警邮件发送逻辑:

import smtplib
import json
from email.mime.text import MIMEText
from google.cloud import secretmanager

def send_alert_email(subject, message_body, recipient_emails):
    # 从Secret Manager读取SMTP配置
    secret_client = secretmanager.SecretManagerServiceClient()
    secret_name = "projects/your-gcp-project-id/secrets/smtp-credentials/versions/latest"
    response = secret_client.access_secret_version(request={"name": secret_name})
    smtp_config = json.loads(response.payload.data.decode("UTF-8"))

    # 构建邮件
    msg = MIMEText(message_body)
    msg["Subject"] = subject
    msg["From"] = smtp_config["smtp_user"]
    msg["To"] = ", ".join(recipient_emails)

    # 发送邮件
    try:
        with smtplib.SMTP(smtp_config["smtp_host"], smtp_config["smtp_port"]) as server:
            server.starttls()
            server.login(smtp_config["smtp_user"], smtp_config["smtp_pass"])
            server.send_message(msg)
        print("告警邮件发送成功")
    except Exception as e:
        print(f"邮件发送失败: {str(e)}")
        # 可将错误日志写入Cloud Logging便于排查

注:若使用Gmail,需开启应用专用密码(不能用普通账号密码)。

3. 在作业中触发告警

根据业务场景(如数据异常、作业失败)调用告警函数:

from pyspark.sql import SparkSession

def main():
    spark = SparkSession.builder.appName("DataprocAlertDemo").getOrCreate()

    try:
        # 示例:读取业务数据
        df = spark.read.parquet("gs://your-bucket/path/to/data")
        
        # 触发条件1:数据量低于阈值
        row_count = df.count()
        if row_count < 1000:
            send_alert_email(
                subject="PySpark作业告警:数据量异常",
                message_body=f"作业执行发现数据行数异常,当前行数:{row_count}",
                recipient_emails=["team-alerts@example.com", "dev@example.com"]
            )
        
        # 其他业务处理逻辑...
        
    except Exception as e:
        # 触发条件2:作业执行异常
        send_alert_email(
            subject="PySpark作业失败告警",
            message_body=f"作业执行失败,错误详情:\n{str(e)}",
            recipient_emails=["on-call@example.com"]
        )
        raise e  # 抛出异常,让Dataproc标记作业状态为失败
    finally:
        spark.stop()

if __name__ == "__main__":
    main()

4. 集群额外配置注意事项

  • 若使用私有集群,需配置VPC防火墙规则,允许集群访问SMTP服务器(如Gmail的smtp.gmail.com:587)。
  • 确保集群服务账号拥有访问目标存储桶、Secret Manager的必要权限。

内容的提问来源于stack exchange,提问作者Mercy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 16:13:28