使用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连接更高效,步骤如下:
确保Spark环境中包含SQL Server JDBC驱动(可以下载对应版本的驱动jar包,放到Spark的
jars目录,或启动时通过--jars参数指定)。编写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
相关产品推荐
相关产品推荐

