借助Azure Databricks将Azure Data Lake Gen2增量变更推送至Azure SQL数据仓库
嗨,我来帮你梳理下具体的实现步骤,还有你关心的低成本集群配置问题~
一、Azure Databricks同步ADLS变更到Azure SQL DW的具体步骤
1. 前置权限与环境准备
- 确保Databricks具备ADLS的访问权限:推荐用**服务主体(Service Principal)**认证,给主体分配ADLS的存储Blob数据 Contributor 角色;也可以用SAS令牌,但服务主体更安全可控。
- 确保Databricks能连接Azure SQL DW:在SQL DW的防火墙规则中允许Databricks集群的IP段(或直接允许Azure服务访问),同时创建具有读写权限的SQL账号(或用服务主体认证SQL DW)。
- 确定ADLS的变更捕获方式:比如按时间戳分区的增量文件、ADLS原生的变更日志(CDC),或者用Azure Event Grid监听文件创建/修改事件,这决定了你后续的增量数据读取逻辑。
2. 挂载ADLS到Databricks
在Databricks笔记本中执行挂载命令,把ADLS容器挂载到Databricks文件系统,方便后续读写:
configs = {"fs.azure.account.auth.type": "OAuth", "fs.azure.account.oauth.provider.type": "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider", "fs.azure.account.oauth2.client.id": "<你的服务主体ID>", "fs.azure.account.oauth2.client.secret": "<你的服务主体密钥>", "fs.azure.account.oauth2.client.endpoint": "https://login.microsoftonline.com/<你的租户ID>/oauth2/token"} # 挂载ADLS容器 dbutils.fs.mount( source = "abfss://<容器名>@<存储账户名>.dfs.core.windows.net/", mount_point = "/mnt/adls-source", extra_configs = configs)
3. 增量数据捕获
根据你的变更方式选择对应逻辑:
- 如果是时间戳分区文件:读取上次同步时间之后的文件,比如:
from pyspark.sql.functions import col # 从控制表读取上次同步的时间(可以存在Delta表或SQL DW的控制表中) last_sync_time = spark.sql("SELECT max(sync_time) FROM control.sync_status").collect()[0][0] # 读取ADLS中的增量数据 incremental_df = spark.read.format("parquet") \ .load("/mnt/adls-source/data") \ .filter(col("modified_time") > last_sync_time)
- 如果是事件驱动:通过Azure Event Grid监听ADLS的文件创建事件,触发Databricks作业,作业直接读取新生成的文件路径。
4. 数据转换与写入SQL DW
完成数据清洗转换后,用Databricks的SQL DW连接器写入目标库,推荐用COPY INTO(高效批量加载)或append模式:
# 配置SQL DW连接参数 sql_dw_url = "jdbc:sqlserver://<SQL DW服务器名>.database.windows.net:1433;database=<数据库名>;encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;loginTimeout=30;" sql_dw_user = "<SQL账号>" sql_dw_password = "<SQL密码>" # 写入SQL DW(append模式) incremental_df.write.format("com.databricks.spark.sqldw") \ .option("url", sql_dw_url) \ .option("dbtable", "target_schema.target_table") \ .option("user", sql_dw_user) \ .option("password", sql_dw_password) \ .mode("append") \ .save() # 如果需要更新/删除,可执行MERGE语句(通过JDBC) merge_query = """ MERGE INTO target_schema.target_table t USING staging_table s ON t.id = s.id WHEN MATCHED THEN UPDATE SET ... WHEN NOT MATCHED THEN INSERT (...) VALUES (...) """ spark.sql(merge_query)
5. 同步状态更新
每次同步完成后,更新同步时间戳到控制表,避免重复同步:
from pyspark.sql.functions import current_timestamp # 更新Delta格式的控制表 spark.sql("INSERT INTO control.sync_status (sync_time) VALUES (current_timestamp())")
二、低成本集群配置:必须借助Databricks作业实现!
是的,要实现集群运行时创建/启动、结束后自动停止/删除,必须依赖Databricks的**作业(Jobs)**功能,手动操作没法自动化做到低成本。这里有两种主流方案:
方案1:使用作业集群(推荐生产环境)
创建作业时选择“新建作业集群”,配置完参数后,开启“作业完成后终止集群”选项:
- 作业触发时自动创建集群,运行结束后立刻终止,完全按需计费,没有闲置成本。
- 可以选择Azure Spot实例(低优先级VM),能节省70%左右的计算成本,只要你的作业能容忍偶尔的中断(作业会自动重试)。
- 触发方式灵活:支持定时触发(比如每天凌晨同步)、事件触发(ADLS有变更就触发)、手动触发。
方案2:交互式集群自动停止(适合测试)
如果是开发测试阶段用笔记本手动运行,可以给交互式集群设置自动停止时间(比如闲置15分钟后停止),但这种方式没法完全自动化,不如作业集群适合生产。
三、作业的具体配置步骤
- 登录Azure Databricks控制台,进入“作业”页面,点击“创建作业”。
- 填写作业名称,选择要运行的笔记本(或JAR/脚本)。
- 在“集群”配置项中,选择“新建作业集群”:
- 设置集群模式为“标准”,选择合适的节点类型(比如
Standard_DS3_v2,根据数据量调整)。 - 勾选“作业完成后终止集群”,确保作业结束后集群自动销毁。
- 可选:启用Spot实例,进一步降低成本。
- 设置集群模式为“标准”,选择合适的节点类型(比如
- 配置触发方式:
- 定时触发:设置Cron表达式(比如
0 0 2 * * ?表示每天凌晨2点运行)。 - 事件触发:通过Azure Event Grid监听ADLS的文件事件,调用Databricks作业API触发运行。
- 定时触发:设置Cron表达式(比如
- 设置重试策略(比如失败后重试3次),避免临时故障导致同步失败。
- 保存作业后,可手动触发测试,或等待定时/事件触发。
内容的提问来源于stack exchange,提问作者Moulshree Suhas
相关产品推荐
相关产品推荐

