如何将Databricks Merge操作结果迁移至Azure SQL Database
将Databricks Merge结果同步至Azure SQL Database
基础方案:全量同步Merge后的Delta表
Merge操作执行完成后,deltaTablePeople对应的Delta表已更新至最新状态,直接读取该表并写入Azure SQL Database即可。
完整代码示例
from delta.tables import * # 原Merge操作逻辑 deltaTablePeople = DeltaTable.forPath(spark, '/tmp/delta/people-10m') deltaTablePeopleUpdates = DeltaTable.forPath(spark, '/tmp/delta/people-10m-updates') dfUpdates = deltaTablePeopleUpdates.toDF() deltaTablePeople.alias('people') \ .merge( dfUpdates.alias('updates'), 'people.id = updates.id' ) \ .whenMatchedUpdate(set = { "id": "updates.id", "firstName": "updates.firstName", "middleName": "updates.middleName", "lastName": "updates.lastName", "gender": "updates.gender", "birthDate": "updates.birthDate", "ssn": "updates.ssn", "salary": "updates.salary" } ) \ .whenNotMatchedInsert(values = { "id": "updates.id", "firstName": "updates.firstName", "middleName": "updates.middleName", "lastName": "updates.lastName", "gender": "updates.gender", "birthDate": "updates.birthDate", "ssn": "updates.ssn", "salary": "updates.salary" } ) \ .execute() # 同步至Azure SQL Database # 替换为你的Azure SQL连接信息 jdbc_url = "jdbc:sqlserver://<SQL服务器名称>.database.windows.net:1433;database=<数据库名称>;" connection_props = { "user": "<用户名>", "password": "<密码>", "driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver" } # 读取Merge后的完整Delta表数据 merged_full_df = deltaTablePeople.toDF() # 写入Azure SQL,根据需求选择模式: # - overwrite: 覆盖目标表所有数据(全量同步常用) # - append: 追加数据(需处理重复键) # - ignore: 表存在则跳过写入 merged_full_df.write.jdbc( url=jdbc_url, table="<目标表名称>", mode="overwrite", properties=connection_props )
进阶方案:增量同步Merge变更行
若无需全量同步,可通过Delta表的变更数据捕获(CDC)功能,仅同步本次Merge中被插入或更新的行,提升同步效率。
前提条件
需先为Delta表启用CDC:
-- 在Databricks SQL或Spark SQL中执行 ALTER TABLE delta.`/tmp/delta/people-10m` SET TBLPROPERTIES (delta.enableChangeDataCapture = true)
增量同步代码示例
from delta.tables import * # 记录Merge前的Delta表版本 before_merge_version = deltaTablePeople.history(1).select("version").collect()[0][0] # 执行原Merge操作(逻辑同基础方案) deltaTablePeople.alias('people') \ .merge( dfUpdates.alias('updates'), 'people.id = updates.id' ) \ .whenMatchedUpdate(set = { "id": "updates.id", "firstName": "updates.firstName", "middleName": "updates.middleName", "lastName": "updates.lastName", "gender": "updates.gender", "birthDate": "updates.birthDate", "ssn": "updates.ssn", "salary": "updates.salary" } ) \ .whenNotMatchedInsert(values = { "id": "updates.id", "firstName": "updates.firstName", "middleName": "updates.middleName", "lastName": "updates.lastName", "gender": "updates.gender", "birthDate": "updates.birthDate", "ssn": "updates.ssn", "salary": "updates.salary" } ) \ .execute() # 读取本次Merge产生的变更数据 change_df = spark.read.format("delta") \ .option("readChangeData", "true") \ .option("startingVersion", before_merge_version + 1) \ .load('/tmp/delta/people-10m') \ .filter("_change_type IN ('insert', 'update_postimage')") # 保留插入行和更新后的镜像行 # 将变更数据写入Azure SQL change_df.write.jdbc( url=jdbc_url, table="<目标表名称>", mode="append", properties=connection_props )
关键注意事项
- JDBC驱动:确保Databricks集群已安装Azure SQL JDBC驱动,可通过集群库管理添加
com.microsoft.sqlserver:mssql-jdbc:12.4.2.jre8(版本可按需调整)。 - 写入模式:全量同步优先用
overwrite,但需注意目标表的索引、约束是否受影响;增量同步用append,需在目标表处理重复键(比如Azure SQL端用MERGE语句)。 - 权限配置:确保Databricks集群能访问Azure SQL,可通过VNet对等连接、防火墙规则或托管身份认证实现。
内容的提问来源于stack exchange,提问作者Patterson
相关产品推荐
相关产品推荐

