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

PySpark通过JDBC调用SQL存储过程报错,求正确实现代码

PySpark通过JDBC调用SQL Server存储过程的正确实现(解决EXEC语法错误)

问题现象

执行PySpark代码调用SQL Server存储过程时,触发语法错误:

Parse error at line: 1, column: 15: Incorrect syntax near 'EXEC'.

错误根源是Spark JDBC的read.jdbc()方法默认将table参数解析为表名,直接传入EXEC语句会被识别为非法的表名格式,导致解析失败。

解决方案1:适合返回单一结果集的存储过程

将EXEC语句包裹为带别名的子查询,让Spark JDBC识别为合法的查询语句:

from pyspark.sql import SparkSession

# 初始化SparkSession
spark = SparkSession.builder.appName("Run SQL Procedure via JDBC").getOrCreate()

# JDBC连接配置
jdbcHostname = "--"
jdbcPort = "--"
jdbcDatabase = "--"
jdbcUsername = "--"
jdbcPassword = "--"

jdbcUrl = f"jdbc:sqlserver://{jdbcHostname}:{jdbcPort};database={jdbcDatabase}"
connectionProperties = {
    "user": jdbcUsername,
    "password": jdbcPassword,
    "driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver"
}

# 存储过程及参数
procedureName = "[jarvis].[uspGetCreditOutputLatestTimeStampForGamification]"
param1 = "DeltaLoad"

# 关键:用括号包裹EXEC语句并指定别名
query = f"(EXEC {procedureName} @param1='{param1}') AS procedure_result"

# 读取存储过程返回的结果集
df = spark.read.jdbc(url=jdbcUrl, table=query, properties=connectionProperties)

# 展示结果
df.show()

解决方案2:适合带输出参数/多结果集的存储过程

直接使用JDBC的CallableStatement API执行存储过程,灵活性更高:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, TimestampType

# 初始化SparkSession
spark = SparkSession.builder.appName("Run Procedure with CallableStatement").getOrCreate()

# JDBC连接配置
jdbcHostname = "--"
jdbcPort = "--"
jdbcDatabase = "--"
jdbcUsername = "--"
jdbcPassword = "--"

jdbcUrl = f"jdbc:sqlserver://{jdbcHostname}:{jdbcPort};database={jdbcDatabase}"

# 存储过程调用格式(?代表参数占位符)
procedure_call = "{call [jarvis].[uspGetCreditOutputLatestTimeStampForGamification](?)}"
param1 = "DeltaLoad"

def execute_stored_procedure():
    conn = None
    stmt = None
    result_set = None
    try:
        # 获取JDBC连接
        conn = spark._sc._jvm.java.sql.DriverManager.getConnection(jdbcUrl, jdbcUsername, jdbcPassword)
        # 创建CallableStatement
        stmt = conn.prepareCall(procedure_call)
        # 设置输入参数(索引从1开始)
        stmt.setString(1, param1)
        
        # 如果有输出参数,需提前注册:例如 stmt.registerOutParameter(2, java.sql.Types.INTEGER)
        
        # 执行存储过程
        has_result = stmt.execute()
        
        # 收集所有结果集数据
        results = []
        while has_result:
            result_set = stmt.getResultSet()
            col_count = result_set.getMetaData().getColumnCount()
            while result_set.next():
                row = [result_set.getObject(i) for i in range(1, col_count + 1)]
                results.append(row)
            has_result = stmt.getMoreResults()
        
        # 定义结果集Schema(需根据存储过程实际返回字段调整)
        schema = StructType([
            StructField("LatestTimeStamp", TimestampType(), nullable=True)
            # 添加其他返回字段...
        ])
        
        # 转换为DataFrame
        return spark.createDataFrame(results, schema)
    
    finally:
        # 关闭资源
        if result_set:
            result_set.close()
        if stmt:
            stmt.close()
        if conn:
            conn.close()

# 执行并展示结果
df = execute_stored_procedure()
df.show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 09:27:02