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

在Azure Databricks中对Azure SQL表执行Upsert/Merge Into的方案咨询

解决Databricks到Azure SQL的UPSERT(更新插入)需求

一、纠正对Merge Into的误解:它支持增量UPSERT,无需全表覆盖

你之前对Merge Into的理解有误——它不是全表覆盖,而是可以精准匹配主键,仅对冲突记录执行更新、新增记录执行插入,完全契合你的需求。在Databricks中通过PySpark结合JDBC连接Azure SQL,就能实现增量Merge Into操作,不会触动整张表的历史数据。

实现示例

方式1:直接通过Spark SQL执行Merge Into

# 假设已处理好今日增量数据DataFrame:incremental_df
# 定义Azure SQL连接参数
jdbc_url = "jdbc:sqlserver://<你的服务器名>.database.windows.net:1433;databaseName=<数据库名>"
connection_properties = {
    "user": "<用户名>",
    "password": "<密码>",
    "driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver"
}

# 将增量数据注册为临时视图
incremental_df.createOrReplaceTempView("daily_updates")

# 执行Merge Into逻辑
spark.sql(f"""
    MERGE INTO delta.`{jdbc_url}`.dbo.<目标表名> AS target
    USING daily_updates AS source
    ON target.<主键列> = source.<主键列>
    WHEN MATCHED THEN
        UPDATE SET 
            target.<列1> = source.<列1>,
            target.<列2> = source.<列2>,
            target.update_time = CURRENT_TIMESTAMP()
    WHEN NOT MATCHED THEN
        INSERT (<主键列>, <列1>, <列2>, update_time)
        VALUES (source.<主键列>, source.<列1>, source.<列2>, CURRENT_TIMESTAMP())
""")

方式2:先写临时表再执行SQL端Merge Into(兼容更多Spark版本)

如果你的Spark版本不直接支持JDBC的Merge Into语法,可先把增量数据写入Azure SQL临时表,再通过JDBC执行SQL端的Merge操作:

# 写入增量数据到临时表
incremental_df.write.jdbc(
    url=jdbc_url,
    table="dbo.staging_daily_updates",
    mode="overwrite",
    properties=connection_properties
)

# 定义Merge逻辑SQL
merge_sql = """
    MERGE INTO dbo.<目标表名> AS target
    USING dbo.staging_daily_updates AS source
    ON target.<主键列> = source.<主键列>
    WHEN MATCHED THEN
        UPDATE SET 
            target.<列1> = source.<列1>,
            target.<列2> = source.<列2>,
            target.update_time = GETDATE()
    WHEN NOT MATCHED THEN
        INSERT (<主键列>, <列1>, <列2>, update_time)
        VALUES (source.<主键列>, source.<列1>, source.<列2>, GETDATE());
    DROP TABLE IF EXISTS dbo.staging_daily_updates;
"""

# 通过JDBC连接执行SQL
conn = spark._sc._gateway.jvm.java.sql.DriverManager.getConnection(
    jdbc_url, connection_properties["user"], connection_properties["password"]
)
stmt = conn.createStatement()
stmt.execute(merge_sql)
stmt.close()
conn.close()

二、从Databricks调用Azure SQL存储过程实现UPSERT

完全可以通过Databricks调用存储过程实现单条/批量的UPSERT逻辑,适合业务规则复杂的场景。

1. 在Azure SQL中创建存储过程

CREATE PROCEDURE dbo.UpsertRecord
    @主键值 INT,
    @列1 VARCHAR(50),
    @列2 INT
AS
BEGIN
    SET NOCOUNT ON;
    IF EXISTS (SELECT 1 FROM dbo.<目标表名> WHERE <主键列> = @主键值)
    BEGIN
        UPDATE dbo.<目标表名>
        SET <列1> = @列1, <列2> = @列2, update_time = GETDATE()
        WHERE <主键列> = @主键值;
    END
    ELSE
    BEGIN
        INSERT INTO dbo.<目标表名> (<主键列>, <列1>, <列2>, update_time)
        VALUES (@主键值, @列1, @列2, GETDATE());
    END
END

2. 在Databricks中调用存储过程

用foreachPartition批量处理,减少JDBC连接次数:

def call_stored_procedure(partition):
    conn = None
    try:
        conn = spark._sc._gateway.jvm.java.sql.DriverManager.getConnection(
            jdbc_url, connection_properties["user"], connection_properties["password"]
        )
        stmt = conn.prepareCall("{call dbo.UpsertRecord(?, ?, ?)}")
        for row in partition:
            stmt.setInt(1, row.<主键列>)
            stmt.setString(2, row.<列1>)
            stmt.setInt(3, row.<列2>)
            stmt.execute()
    finally:
        if conn is not None:
            conn.close()

# 对增量DataFrame执行批量调用
incremental_df.foreachPartition(call_stored_procedure)

三、性能优化建议

  • 增量数据量较大时,优先选择临时表+SQL端Merge Into的方式,比逐条调用存储过程效率更高。
  • 增量数据量小、业务规则复杂时,存储过程的方式更灵活。
  • 确保Azure SQL的主键列已创建索引,避免Merge或更新时全表扫描,提升执行速度。

内容的提问来源于stack exchange,提问作者Sivaani

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 07:56:21