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

如何将Databricks Delta文件中不存在的记录导入Azure SQL已有表

仅导入Delta表中不存在于Azure SQL的记录的解决方案

下面提供几种可行的实现方式,解决你遇到的重复导入问题:

方法一:Spark左反连接过滤新记录

先读取Azure SQL目标表的唯一标识字段,通过左反连接(left_anti)筛选出Delta表中未存在于目标表的记录,再执行追加写入。

// 读取Delta源数据
val deltaDF = spark.read.format("delta").load("dbfs:/user/hive/warehouse/hourly_tables_JMA/")

// 读取Azure SQL目标表的唯一标识列(替换为你表的主键/唯一键,比如id)
val sqlExistingIdsDF = spark.read
  .format("jdbc")
  .option("url", "jdbc:sqlserver://qdenaz.database.windows.net:1433;databaseName=database")
  .option("dbtable", "[DWH].HOURLY_TABLES_JMA")
  .option("user", "Admin")
  .option("password", "****")
  .load()
  .select("id") // 仅读取标识列,减少数据传输量

// 筛选出Delta中不存在于SQL表的记录
val newRecordsDF = deltaDF.join(sqlExistingIdsDF, deltaDF("id") === sqlExistingIdsDF("id"), "left_anti")

// 写入新记录到目标表
newRecordsDF.write 
  .format("jdbc")
  .option("url", "jdbc:sqlserver://qdenaz.database.windows.net:1433;databaseName=database")
  .option("dbtable", "[DWH].HOURLY_TABLES_JMA")
  .mode("append")   
  .option("user", "Admin")
  .option("password", "****")
  .save()

注意:如果目标表数据量极大,建议添加过滤条件(比如按时间范围读取SQL表的近期记录),或者结合Delta表的增量读取特性,避免全量读取带来的性能问题。

方法二:使用Azure SQL的MERGE语句

将Delta数据写入SQL临时表,再通过MERGE语句仅插入不存在的记录,这种方式无需将整个目标表加载到Spark,适合大数据量场景。

// 读取Delta源数据
val deltaDF = spark.read.format("delta").load("dbfs:/user/hive/warehouse/hourly_tables_JMA/")

// 将Delta数据写入Azure SQL临时表
deltaDF.write 
  .format("jdbc")
  .option("url", "jdbc:sqlserver://qdenaz.database.windows.net:1433;databaseName=database")
  .option("dbtable", "#temp_hourly_jma") // 会话级临时表,关闭后自动删除
  .mode("overwrite")   
  .option("user", "Admin")
  .option("password", "****")
  .save()

// 构造MERGE语句,替换成你表的主键匹配条件和字段列表
val mergeSql = """
MERGE INTO [DWH].HOURLY_TABLES_JMA target
USING #temp_hourly_jma source
ON target.id = source.id -- 主键/唯一匹配条件
WHEN NOT MATCHED THEN
INSERT (col1, col2, col3, ...) -- 目标表所有字段
VALUES (source.col1, source.col2, source.col3, ...)
"""

// 执行MERGE语句
val conn = java.sql.DriverManager.getConnection(
  "jdbc:sqlserver://qdenaz.database.windows.net:1433;databaseName=database",
  "Admin",
  "****"
)
val stmt = conn.createStatement()
stmt.execute(mergeSql)
stmt.close()
conn.close()

方法三:基于Delta Lake的CDC增量同步

如果你的Delta表开启了变更数据捕获(CDC),可以直接读取新增的插入记录,避免重复处理已有数据。

第一步:开启Delta表CDC(仅需执行一次)

spark.sql("ALTER TABLE hourly_tables_JMA SET TBLPROPERTIES (delta.enableChangeDataCapture = true)")

第二步:读取增量数据并写入SQL

// 维护上次同步的版本号(可以存在数据库或配置文件中)
val lastSyncedVersion = 123 // 替换为实际的上次同步版本

// 读取上次版本之后的新增插入记录
val incrementalDF = spark.read.format("delta")
  .option("readChangeFeed", "true")
  .option("startingVersion", lastSyncedVersion)
  .load("dbfs:/user/hive/warehouse/hourly_tables_JMA/")
  .filter("_change_type = 'insert'") // 仅保留新增记录

// 写入Azure SQL
incrementalDF.write 
  .format("jdbc")
  .option("url", "jdbc:sqlserver://qdenaz.database.windows.net:1433;databaseName=database")
  .option("dbtable", "[DWH].HOURLY_TABLES_JMA")
  .mode("append")   
  .option("user", "Admin")
  .option("password", "****")
  .save()

// 更新lastSyncedVersion为当前Delta表的最新版本,下次同步使用
val currentVersion = spark.sql("DESCRIBE HISTORY hourly_tables_JMA").select("version").head().getLong(0)

内容的提问来源于stack exchange,提问作者Jose Maria Acevedo Reinoso

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 02:10:55