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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 02:55:33