如何从Azure Databricks对Azure SQL Server表执行Upsert操作?
从Azure Databricks对Azure SQL数据库执行Upsert(当日数据替换/插入)
1. 准备当日处理后的数据集
先确保你的DataFrame仅包含当日生成的新记录,提前过滤日期条件:
# 假设df是处理完成的全量数据,process_date为数据所属日期字段 from pyspark.sql.functions import current_date filtered_df = df.filter(df.process_date == current_date())
2. 配置Azure SQL连接信息
使用Databricks密钥管理存储敏感信息,避免硬编码密码:
server_name = "tcp:你的SQL服务器名.database.windows.net" database_name = "你的数据库名" jdbc_url = f"jdbc:sqlserver://{server_name};databaseName={database_name};encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;loginTimeout=30;" connection_properties = { "user" : "SQL账号", "password" : dbutils.secrets.get(scope="你的密钥范围", key="sql-password"), "driver" : "com.microsoft.sqlserver.jdbc.SQLServerDriver" }
3. 执行Upsert的两种实现方式
方式一:临时表+MERGE INTO(推荐,适配大数据量)
先将当日数据写入Azure SQL的全局临时表,再通过MERGE语句完成替换/插入:
# 将当日数据写入全局临时表 filtered_df.write.jdbc( url=jdbc_url, table="##temp_daily_records", mode="overwrite", properties=connection_properties ) # 定义MERGE逻辑SQL语句 merge_sql = """ MERGE INTO 目标表名 t USING ##temp_daily_records s -- 匹配条件:当日日期+业务主键,确保仅更新当日的同一条记录 ON t.process_date = s.process_date AND t.业务主键字段 = s.业务主键字段 WHEN MATCHED THEN UPDATE SET t.字段1 = s.字段1, t.字段2 = s.字段2, -- 列出所有需要更新的字段 t.last_updated = GETDATE() WHEN NOT MATCHED THEN INSERT (process_date, 业务主键字段, 字段1, 字段2, ...) VALUES (s.process_date, s.业务主键字段, s.字段1, s.字段2, ...); -- 清理临时表 DROP TABLE IF EXISTS ##temp_daily_records; """ # 执行MERGE语句 spark.read.jdbc(url=jdbc_url, table=f"({merge_sql}) AS merge_result", properties=connection_properties)
方式二:直接JDBC执行MERGE(仅适合小数据量)
如果数据量很小,可以直接构造MERGE的VALUES子句,但需注意SQL注入风险:
# 定义MERGE模板 merge_sql = """ MERGE INTO 目标表名 t USING (VALUES {value_list}) s (process_date, 业务主键字段, 字段1, 字段2) ON t.process_date = s.process_date AND t.业务主键字段 = s.业务主键字段 WHEN MATCHED THEN UPDATE SET 字段1 = s.字段1, 字段2 = s.字段2 WHEN NOT MATCHED THEN INSERT (process_date, 业务主键字段, 字段1, 字段2) VALUES (s.process_date, s.业务主键字段, s.字段1, s.字段2); """ # 构造VALUES列表,处理字符串转义 value_list = [] for row in filtered_df.collect(): # 对字符串字段转义单引号 col1 = f"'{row.字段1.replace(\"'\", \"''\")}'" if isinstance(row.字段1, str) else row.字段1 col2 = f"'{row.字段2.replace(\"'\", \"''\")}'" if isinstance(row.字段2, str) else row.字段2 value_list.append(f"('{row.process_date}', {row.业务主键字段}, {col1}, {col2})") final_merge_sql = merge_sql.format(value_list=", ".join(value_list)) # 执行SQL spark.read.jdbc(url=jdbc_url, table=f"({final_merge_sql}) AS merge_result", properties=connection_properties)
关键注意事项
- 精准匹配:
ON子句必须包含当日日期+业务主键,避免误更新非当日数据。 - 性能优化:大数据量优先用临时表方案;给目标表的
process_date和主键字段建立索引,提升MERGE执行效率。 - 权限配置:Databricks服务账号需拥有Azure SQL的
INSERT、UPDATE、CREATE TABLE权限。 - 原子性:MERGE语句本身是原子操作,确保更新/插入要么全部成功,要么回滚。
内容的提问来源于stack exchange,提问作者Sivaani N
相关产品推荐
相关产品推荐

