如何通过Azure Databricks避免向Azure SQL DB插入重复数据
避免ADLS数据导入Azure SQL DB mysalesorder表时产生重复行的方案
方案1:借助Azure SQL的MERGE语句(推荐)
先在Azure SQL的mysalesorder表上定义唯一键/主键(比如以order_id这类业务唯一标识作为主键),然后通过Databricks执行MERGE语句,实现"存在则更新,不存在则插入"的逻辑,从根源避免重复。
示例代码:
# 假设你的数据源DataFrame是sales_df,包含order_id等核心字段 # 先将DataFrame写入Azure SQL的临时表 sales_df.write \ .format("jdbc") \ .option("url", "jdbc:sqlserver://<server-name>.database.windows.net:1433;databaseName=<db-name>") \ .option("dbtable", "#temp_sales") \ .option("user", "<username>") \ .option("password", "<password>") \ .mode("overwrite") \ .save() # 编写MERGE合并逻辑 merge_query = """ MERGE INTO mysalesorder AS target USING #temp_sales AS source ON target.order_id = source.order_id WHEN MATCHED THEN UPDATE SET target.order_date = source.order_date, target.customer_id = source.customer_id, -- 按需添加其他需要更新的字段 WHEN NOT MATCHED THEN INSERT (order_id, order_date, customer_id, ...) VALUES (source.order_id, source.order_date, source.customer_id, ...); """ # 通过JDBC执行MERGE语句 spark.read.jdbc( url="jdbc:sqlserver://<server-name>.database.windows.net:1433;databaseName=<db-name>", table=f"({merge_query}) AS merge_result", properties={"user": "<username>", "password": "<password>"} ).show()
方案2:在Databricks端先去重再导入
如果不想在SQL端处理逻辑,可以直接在Databricks的DataFrame层面,基于业务唯一键完成去重后再插入正式表。
示例代码:
# 以order_id为唯一标识,对DataFrame去重 deduplicated_df = sales_df.dropDuplicates(["order_id"]) # 写入正式表,使用append模式 deduplicated_df.write \ .format("jdbc") \ .option("url", "jdbc:sqlserver://<server-name>.database.windows.net:1433;databaseName=<db-name>") \ .option("dbtable", "mysalesorder") \ .option("user", "<username>") \ .option("password", "<password>") \ .mode("append") \ .save()
注意:这种方式依赖唯一键的准确性,适合数据源本身重复量少的场景,否则可能出现漏判。
方案3:导入前比对已有数据
先从Azure SQL读取mysalesorder表的唯一键集合,再在Databricks中过滤掉已存在的数据,只插入全新的行。
示例代码:
# 读取SQL表中已存在的order_id集合 existing_order_ids = spark.read \ .format("jdbc") \ .option("url", "jdbc:sqlserver://<server-name>.database.windows.net:1433;databaseName=<db-name>") \ .option("dbtable", "SELECT order_id FROM mysalesorder") \ .option("user", "<username>") \ .option("password", "<password>") \ .load() # 过滤出DataFrame中未存在的新行 new_sales_df = sales_df.join(existing_order_ids, on="order_id", how="left_anti") # 插入新行到正式表 new_sales_df.write \ .format("jdbc") \ .option("url", "jdbc:sqlserver://<server-name>.database.windows.net:1433;databaseName=<db-name>") \ .option("dbtable", "mysalesorder") \ .option("user", "<username>") \ .option("password", "<password>") \ .mode("append") \ .save()
注意:如果SQL表数据量极大,读取唯一键集合会占用较多资源,更适合中小数据量场景。
内容的提问来源于stack exchange,提问作者user8205502
相关产品推荐
相关产品推荐

