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

如何通过Synapse的Spark池将DataFrame数据追加到专用SQL池

Azure Synapse Spark池DataFrame追加写入专用SQL池操作指南

前置校验

  • 确认Spark池和目标专用SQL池归属同一Synapse工作区,你的操作账号同时持有两个资源的数据写入权限
  • 提前在专用SQL池中创建好目标表,表字段的类型、顺序需和待写入DataFrame完全匹配,避免隐式转换报错
  • 无需额外安装连接器依赖,Synapse运行环境已预装适配的SQL DW连接器

PySpark实现代码

核心使用内置的com.databricks.spark.sqldw连接器,指定append写入模式即可实现追加写入,示例代码如下:

# 1. 构造或读取待写入的DataFrame(此处为示例测试数据,可替换为你自己的数据源读取逻辑)
test_df = spark.createDataFrame(
    [
        (1, "张三", 28, "技术部"),
        (2, "李四", 32, "产品部"),
        (3, "王五", 25, "运营部")
    ],
    ["id", "name", "age", "department"]
)

# 2. 填写连接配置参数,替换<>占位符为你的实际资源信息
jdbc_url = "jdbc:sqlserver://<你的Synapse工作区名>.sql.azuresynapse.net:1433;database=<专用SQL池名称>;encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.sql.azuresynapse.net;loginTimeout=30;"
sql_user = "<SQL池管理员账号>"
sql_pwd = "<SQL池管理员密码>"
target_table = "<目标表完整名,例:dbo.staff_info>"
temp_dir = "abfss://<存储容器名>@<关联存储账户名>.dfs.core.windows.net/<临时目录自定义路径>" # 用于暂存写入过程中的中间数据,Synapse会自动清理过期文件

# 3. 执行追加写入
test_df.write \
  .format("com.databricks.spark.sqldw") \
  .option("url", jdbc_url) \
  .option("user", sql_user) \
  .option("password", sql_pwd) \
  .option("dbtable", target_table) \
  .option("tempDir", temp_dir) \
  .option("forwardSparkAzureStorageCredentials", "true") \
  .mode("append") \
  .save()

优化与注意事项

  • 敏感信息不要硬编码到代码中,可将账号密码存入工作区关联的密钥保管库,运行时通过mssparkutils.credentials.getSecret("密钥保管库名", "密钥名")读取
  • 如需使用AAD托管身份认证,可删除user、password配置项,新增option("useAzureADIdentity", "true")即可,前提是Spark池的托管身份已被授予专用SQL池的写入权限
  • 大数据量写入前可对DataFrame做重分区操作,提升并行写入效率,可通过option("maxErrors", 允许的最大错误行数)配置容错阈值,少量脏数据不会导致整个写入任务失败
  • 若目标表存在主键、唯一约束等限制,需提前校验待写入数据不会违反约束,避免写入报错

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 01:24:03