如何在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
相关产品推荐
相关产品推荐

