在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
相关产品推荐
相关产品推荐

