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

如何在Databricks中实现传入文件审计及零字节校验的PySpark代码

Databricks PySpark 零字节文件校验实现方案

实现逻辑说明

  • 直接通过Databricks内置的dbutils工具读取文件元信息判断大小,相比读取文件内容判断效率更高,适合大文件场景
  • 校验到文件为0字节时触发邮件告警流程,同时终止当前作业;校验通过则继续执行后续业务逻辑
  • 兼容DBFS路径、挂载的云存储(S3/ADLS/GCS等)路径的文件校验

完整实现代码

import smtplib
from email.mime.text import MIMEText
from email.header import Header

# -------------------------- 可修改配置项 --------------------------
# 待校验的目标文件路径
TARGET_FILE_PATH = "/mnt/your_input_directory/your_target_file.csv"
# 邮件服务配置,可根据企业实际SMTP服务调整
SMTP_SERVER = "smtp.your-company.com"
SMTP_PORT = 587
SENDER_EMAIL = "databricks_alert@your-company.com"
SENDER_AUTH_CODE = "your_smtp_auth_code"
ALERT_RECEIVERS = ["data_ops@your-company.com", "data_owner@your-company.com"]
# -------------------------------------------------------------------

def send_zero_byte_alert(file_path: str) -> None:
    """发送零字节文件告警邮件"""
    mail_content = f"""
    【Databricks数据审计告警】
    检测到传入文件为零字节空文件,当前作业已终止。
    问题文件路径:{file_path}
    所属作业ID:{dbutils.notebook.entry_point.getDbutils().notebook().getContext().jobId().getOrElse("手动运行作业")}
    触发时间:{dbutils.notebook.entry_point.getDbutils().notebook().getContext().startTime().get()}
    """
    msg = MIMEText(mail_content, "plain", "utf-8")
    msg["From"] = Header("Databricks数据审计平台", "utf-8")
    msg["To"] = Header("数据运维及业务负责人", "utf-8")
    msg["Subject"] = Header("零字节文件告警,作业已终止", "utf-8")

    try:
        smtp_client = smtplib.SMTP(SMTP_SERVER, SMTP_PORT)
        smtp_client.starttls()
        smtp_client.login(SENDER_EMAIL, SENDER_AUTH_CODE)
        smtp_client.sendmail(SENDER_EMAIL, ALERT_RECEIVERS, msg.as_string())
        smtp_client.quit()
        print("告警邮件已发送至相关负责人")
    except Exception as e:
        print(f"告警邮件发送失败,错误信息:{str(e)}")
        raise e

def is_zero_byte_file(file_path: str) -> bool:
    """校验目标文件是否为零字节空文件,返回True代表为空文件"""
    try:
        file_meta = dbutils.fs.ls(file_path)[0]
        return file_meta.size == 0
    except Exception as e:
        print(f"读取文件元信息失败,错误信息:{str(e)}")
        raise e

# 主校验流程
if __name__ == "__main__":
    empty_flag = is_zero_byte_file(TARGET_FILE_PATH)
    if empty_flag:
        print(f"校验不通过:文件{TARGET_FILE_PATH}为零字节空文件,终止作业")
        send_zero_byte_alert(TARGET_FILE_PATH)
        # 终止当前notebook运行,返回失败状态
        dbutils.notebook.exit("零字节文件校验失败,作业终止")
    else:
        print(f"校验通过:文件{TARGET_FILE_PATH}存在有效数据,继续后续处理")
        # 此处写入后续业务处理代码
        # 示例:df = spark.read.csv(TARGET_FILE_PATH, header=True)
        # ...

适配场景说明

  • 如果需要校验目录下所有文件,可遍历dbutils.fs.ls(目标目录路径)返回的所有文件对象逐个校验即可
  • 若企业使用云厂商托管的邮件服务(如AWS SES、Azure SendGrid),可替换邮件发送逻辑为对应云服务的调用代码
  • 单文件大小超过TB级的场景也可正常使用,该校验逻辑不会读取文件实际内容,仅拉取元数据,性能开销极低

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 18:57:03