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

如何从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 07:32:25