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

使用Python连接Azure SQL Server并通过PySpark查询时遇连接错误

解决Azure SQL Server连接及PySpark查询问题

先修复pymssql连接错误

你的代码存在几个关键问题,对应修正方案如下:

  • 用户名格式错误:Azure SQL Server不支持域\用户名的格式,需改为用户名@服务器前缀,比如把Mfc\SVC_Account改成SVC_Account@MFCGaZURESERVERHOST(服务器前缀是你主机名MFCGaZURESERVERHOST.database.windows.net的前半部分)。
  • 转义字符问题:原代码里的\是Python转义符,若坚持用域格式(仅适用于Azure AD认证场景,pymssql支持有限),需改成双反斜杠\\或原始字符串r"Mfc\SVC_Account"。
  • 防火墙未开放:登录Azure门户,找到你的SQL Server,在「防火墙和虚拟网络」中添加当前客户端IP,或临时允许所有IP(生产环境禁用)。

修正后的pymssql代码:

import pymssql

# 注意用户名格式改为 用户名@服务器前缀
conn = pymssql.connect(
    server="MFCGaZURESERVERHOST.database.windows.net",
    user="SVC_Account@MFCGaZURESERVERHOST",
    password="rYa9NqsUpuN2kCm",
    database="jp-database_Name"
)

cursor = conn.cursor()
cursor.execute("DELETE from SELECTIONTAIL where 1=1")
conn.commit()
conn.close()  # 记得关闭连接

如果还是报错,建议改用pyodbc(对Azure SQL支持更完善),示例代码:

import pyodbc

conn_str = (
    "DRIVER={ODBC Driver 18 for SQL Server};"
    "SERVER=MFCGaZURESERVERHOST.database.windows.net;"
    "DATABASE=jp-database_Name;"
    "UID=SVC_Account@MFCGaZURESERVERHOST;"
    "PWD=rYa9NqsUpuN2kCm;"
    "Encrypt=yes;"
    "TrustServerCertificate=no;"
    "Connection Timeout=30;"
)
conn = pyodbc.connect(conn_str)
# 后续操作和pymssql一致

用PySpark连接Azure SQL执行查询

既然你的目标是用PySpark处理大型数据集,直接用Spark的JDBC连接更高效,步骤如下:

  1. 确保Spark环境中包含SQL Server JDBC驱动(可以下载对应版本的驱动jar包,放到Spark的jars目录,或启动时通过--jars参数指定)。

  2. 编写PySpark代码:

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("Azure SQL Query") \
    .config("spark.driver.extraClassPath", "/path/to/mssql-jdbc-xxx.jar")  # 替换为你的JDBC驱动路径
    .getOrCreate()

# 连接参数
server_name = "MFCGaZURESERVERHOST.database.windows.net"
database_name = "jp-database_Name"
url = f"jdbc:sqlserver://{server_name}:1433;databaseName={database_name};encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;loginTimeout=30;"
table_name = "SELECTIONTAIL"
properties = {
    "user": "SVC_Account@MFCGaZURESERVERHOST",
    "password": "rYa9NqsUpuN2kCm",
    "driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver"
}

# 读取数据(示例)
df = spark.read.jdbc(url=url, table=table_name, properties=properties)
df.show()

# 执行删除操作(两种方式可选)
# 方式1:通过Spark SQL执行
df.createOrReplaceTempView("selection_tail")
spark.sql("DELETE FROM selection_tail WHERE 1=1")

# 方式2:用JDBC连接直接执行(适合批量操作)
connection = spark._sc._gateway.jvm.java.sql.DriverManager.getConnection(url, properties["user"], properties["password"])
statement = connection.createStatement()
statement.execute("DELETE from SELECTIONTAIL where 1=1")
connection.commit()
connection.close()

spark.stop()

注意:

  • 执行删除等写操作时,Spark的JDBC写入性能取决于分区设置,若数据集很大,建议分批次处理。
  • 生产环境避免用WHERE 1=1全表删除,确认逻辑后再执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 14:35:24