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

如何通过Service Principal认证高效写入Azure Synapse dedicated SQL池(含Databricks代码需求)

问题1:通过Service Principal认证从Synapse写入Azure Synapse Dedicated SQL池

前置准备

  • 注册Azure AD服务主体(SP),记录其tenantId、clientId、clientSecret
  • 给SP分配Synapse工作区权限:至少需要Synapse Contributor或更细粒度权限,确保能访问目标SQL池
  • 在目标Dedicated SQL池中给SP分配数据库权限:添加为db_datawriter角色,或直接授予INSERT、UPDATE等写入权限

Synapse Notebook PySpark代码示例

# SP认证配置
tenant_id = "<你的租户ID>"
client_id = "<你的SP客户端ID>"
client_secret = "<你的SP客户端密钥>"
sql_pool_server = "<SQL池服务器名>.database.windows.net"
sql_pool_database = "<目标数据库名>"
target_table = "<目标表名>"

# 设置Spark OAuth认证参数
spark.conf.set(f"spark.sql.azure.synapse.service.principal.client.id", client_id)
spark.conf.set(f"spark.sql.azure.synapse.service.principal.client.secret", client_secret)
spark.conf.set(f"spark.sql.azure.synapse.service.principal.tenant.id", tenant_id)

# 读取源数据(示例为生成测试DataFrame)
source_df = spark.range(1000).selectExpr("id as user_id", "concat('user_', id) as user_name")

# 写入Dedicated SQL池
source_df.write \
    .format("com.microsoft.sqlserver.jdbc.spark") \
    .mode("append") \
    .option("url", f"jdbc:sqlserver://{sql_pool_server}:1433;database={sql_pool_database};encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;loginTimeout=30;") \
    .option("dbtable", target_table) \
    .option("authentication", "ActiveDirectoryServicePrincipal") \
    .option("user", client_id) \
    .option("password", client_secret) \
    .save()

问题2:Databricks/PySpark通过SP认证高性能写入Dedicated SQL池

JDBC速度慢的核心原因是单批次小批量写入,推荐使用Azure Synapse Connector for Spark(底层基于COPY INTO命令,支持并行批量写入),且支持SP OAuth认证,无需存储账户密钥。

前置准备

  • 同问题1的SP权限配置,额外确保SP对SQL池有CONTROL权限(或至少能执行COPY INTO命令)
  • 在Databricks集群中安装Synapse Connector依赖:通过集群库添加com.microsoft.azure:synapsespark_2.12:1.0.0(版本根据集群Scala版本调整)

PySpark代码示例(SP认证)

# SP认证配置
tenant_id = "<你的租户ID>"
client_id = "<你的SP客户端ID>"
client_secret = "<你的SP客户端密钥>"
sql_pool_server = "<SQL池服务器名>.database.windows.net"
sql_pool_database = "<目标数据库名>"
target_table = "<目标表名>"

# 设置Spark配置
spark.conf.set("spark.sql.azure.synapse.service.principal.client.id", client_id)
spark.conf.set("spark.sql.azure.synapse.service.principal.client.secret", client_secret)
spark.conf.set("spark.sql.azure.synapse.service.principal.tenant.id", tenant_id)

# 读取源数据(示例为读取Parquet文件)
source_df = spark.read.parquet("/path/to/source/data")

# 高性能写入配置:开启并行写入
source_df.write \
    .format("com.microsoft.azure.synapse") \
    .mode("append") \
    .option("url", f"jdbc:sqlserver://{sql_pool_server}:1433;database={sql_pool_database};") \
    .option("dbtable", target_table) \
    .option("authentication", "ActiveDirectoryServicePrincipal") \
    .option("user", client_id) \
    .option("password", client_secret) \
    .option("numPartitions", "8")  # 并行分区数,根据集群核心数调整
    .option("batchSize", "100000")  # 每个批次写入行数
    .save()

Scala代码示例(SP认证)

import org.apache.spark.sql.SaveMode

// SP认证配置
val tenantId = "<你的租户ID>"
val clientId = "<你的SP客户端ID>"
val clientSecret = "<你的SP客户端密钥>"
val sqlPoolServer = "<SQL池服务器名>.database.windows.net"
val sqlPoolDatabase = "<目标数据库名>"
val targetTable = "<目标表名>"

// 设置Spark配置
spark.conf.set("spark.sql.azure.synapse.service.principal.client.id", clientId)
spark.conf.set("spark.sql.azure.synapse.service.principal.client.secret", clientSecret)
spark.conf.set("spark.sql.azure.synapse.service.principal.tenant.id", tenantId)

// 读取源数据(示例为读取Parquet文件)
val sourceDF = spark.read.parquet("/path/to/source/data")

// 高性能写入
sourceDF.write
  .format("com.microsoft.azure.synapse")
  .mode(SaveMode.Append)
  .option("url", s"jdbc:sqlserver://$sqlPoolServer:1433;database=$sqlPoolDatabase;")
  .option("dbtable", targetTable)
  .option("authentication", "ActiveDirectoryServicePrincipal")
  .option("user", clientId)
  .option("password", clientSecret)
  .option("numPartitions", "8")
  .option("batchSize", "100000")
  .save()

性能优化要点

  • 并行分区数:根据Databricks集群核心数调整,建议每个分区对应1-2个核心,避免资源浪费或过载
  • 批次大小:根据表结构复杂度调整,一般5-10万行/批次较合适
  • 数据预处理:写入前过滤、聚合数据,减少写入量
  • 表结构优化:目标表使用列存储索引,避免频繁Schema变更

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 09:37:50