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

