如何在Databricks中创建每月自动运行的Delta Live Tables数据管道
在Databricks门户内设置Delta Live Tables每月自动同步多数据源的方案
一、先搭建Delta Live Tables(DLT)数据管道
- 登录Databricks后,进入工作流 > Delta Live Tables,点击创建管道
- 填写基础配置:
- 给管道命名,选择集群模式(推荐用服务器less,无需手动维护集群)
- 指定存储位置(用来存放管道的元数据和运行日志)
- 编写/上传DLT处理脚本:
脚本要包含从多数据源拉取、处理数据的逻辑,用DLT的装饰器定义表。举个实际示例:import dlt from pyspark.sql.functions import current_timestamp # 拉取MySQL数据源的原始订单数据 @dlt.table(name="raw_order_data", comment="每月同步的原始订单数据") def pull_raw_orders(): return ( spark.read .format("jdbc") .option("url", "jdbc:mysql://db-host:3306/order_db") .option("dbtable", "orders") .option("user", "db_user") .option("password", "db_pwd") .load() .withColumn("ingest_time", current_timestamp()) ) # 拉取对象存储上的CSV日志数据 @dlt.table(name="raw_log_data", comment="每月同步的用户行为日志") def pull_raw_logs(): return ( spark.read .csv("s3://your-bucket/logs/monthly/", header=True, inferSchema=True) .withColumn("ingest_time", current_timestamp()) ) # 清洗合并后的业务表 @dlt.table(name="merged_order_logs", comment="清洗合并后的订单-日志关联数据") def merge_data(): orders = dlt.read("raw_order_data") logs = dlt.read("raw_log_data") return orders.join(logs, on="user_id", how="inner").filter("order_status = 'completed'") - 点击创建后,先手动运行一次管道,确认数据源连接、数据处理逻辑都正常,避免调度后出问题
二、配置每月自动调度
- 回到你的DLT管道页面,点击右上角的调度按钮
- 在调度面板里完成以下设置:
- 打开已调度开关
- 调度类型选cron(灵活适配每月周期)
- 输入cron表达式:比如
0 0 1 * *代表每月1号凌晨0点执行;如果要每月10号下午2点,就写0 14 10 * *(cron格式为分 时 日 月 周) - 选择对应的时区,确保执行时间符合业务需求
- 可选设置重试策略,防止临时故障导致任务失败
- 点击保存,调度就正式生效了
三、验证与维护
- 在管道的历史标签页,可以看到未来的调度任务列表
- 可以手动触发一次调度任务,验证整个流程是否正常运行
- 后续要修改调度周期或管道逻辑,直接回到对应页面调整即可
注意要点
- 确保你的数据源和Databricks集群网络互通(比如私有数据库需要配置VPC peering,或者用Databricks的私有访问通道)
- 服务器less集群会自动根据任务负载调整资源,适合每月一次的周期性任务
- 可以通过DLT的监控标签查看数据质量报告、运行状态,快速排查问题
内容的提问来源于stack exchange,提问作者billherd67
相关产品推荐
相关产品推荐

