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

AWS Glue执行SQL Server存储过程仅返回true,无法获取结果集求助

在AWS Glue中调用SQL Server存储过程无法获取结果集的解决方法

我在AWS Glue中已成功连接SQL Server并执行存储过程MNA.dbo.zz_MNAvArticulosListar,但仅返回布尔值true,无法获取存储过程返回的结果集。此前可在DBeaver中正常调用该存储过程,原代码如下:

import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job

args = getResolvedOptions(sys.argv, ['TempDir','JOB_NAME'])

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

from py4j.java_gateway import java_import
java_import(sc._gateway.jvm,"java.sql.Connection")
java_import(sc._gateway.jvm,"java.sql.DatabaseMetaData")
java_import(sc._gateway.jvm,"java.sql.DriverManager")
java_import(sc._gateway.jvm,"java.sql.SQLException")

print('Trying to connect to DB')
source_jdbc_conf = glueContext.extract_jdbc_conf('sgc_con')
conn = sc._gateway.jvm.DriverManager.getConnection(source_jdbc_conf.get('url'), source_jdbc_conf.get('user'), source_jdbc_conf.get('password'))
print('Trying to connect to DB success!')
print(conn.getMetaData())
print('prepareCall')
statement = "EXEC MNA.dbo.zz_MNAvArticulosListar"
exec_statement = conn.prepareCall(statement)
print('execute')
rs = exec_statement.execute()
print(rs) #true
exec_statement.close()

问题原因

exec_statement.execute()的返回值仅表示当前存在结果集(true)或不存在(false),并非结果本身。要获取实际数据,需要通过CallableStatement的API提取结果集,同时需处理存储过程可能返回的多个结果集或更新计数。

解决方法

方法1:优化原生Java连接代码

手动处理结果集循环,提取数据并转换为Spark DataFrame:

import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from py4j.java_gateway import java_import

args = getResolvedOptions(sys.argv, ['TempDir','JOB_NAME'])

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

# 导入Java SQL依赖类
java_import(sc._gateway.jvm,"java.sql.Connection")
java_import(sc._gateway.jvm,"java.sql.DatabaseMetaData")
java_import(sc._gateway.jvm,"java.sql.DriverManager")
java_import(sc._gateway.jvm,"java.sql.SQLException")

print('Connecting to DB')
source_jdbc_conf = glueContext.extract_jdbc_conf('sgc_con')
conn = sc._gateway.jvm.DriverManager.getConnection(source_jdbc_conf.get('url'), source_jdbc_conf.get('user'), source_jdbc_conf.get('password'))
print('DB connection success!')

statement = "EXEC MNA.dbo.zz_MNAvArticulosListar"
exec_statement = conn.prepareCall(statement)

print('Executing stored procedure')
has_result = exec_statement.execute()

# 循环处理所有结果集(存储过程可能返回多个结果)
while has_result or exec_statement.getUpdateCount() != -1:
    if has_result:
        # 获取结果集并转换为Spark DataFrame
        result_set = exec_statement.getResultSet()
        df = spark.read.jdbc(
            url=source_jdbc_conf.get('url'),
            table=f"({statement}) AS temp_result",
            properties={
                "user": source_jdbc_conf.get('user'),
                "password": source_jdbc_conf.get('password'),
                "driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver"
            }
        )
        # 打印结果或执行后续处理逻辑
        df.show()
    # 检查是否有更多结果集
    has_result = exec_statement.getMoreResults()

# 关闭资源
exec_statement.close()
conn.close()
job.commit()

方法2:直接使用Spark JDBC(推荐)

无需手动管理Java连接,Spark会自动处理存储过程的结果集,更适配Glue的Spark环境:

import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job

args = getResolvedOptions(sys.argv, ['TempDir','JOB_NAME'])

sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)

source_jdbc_conf = glueContext.extract_jdbc_conf('sgc_con')

# 直接通过Spark JDBC执行存储过程,结果转为DataFrame
df = spark.read.jdbc(
    url=source_jdbc_conf.get('url'),
    table="(EXEC MNA.dbo.zz_MNAvArticulosListar) AS result_set",
    properties={
        "user": source_jdbc_conf.get('user'),
        "password": source_jdbc_conf.get('password'),
        "driver": "com.microsoft.sqlserver.jdbc.SQLServerDriver"
    }
)

# 查看结果
df.show()

# 可选:将结果写入S3或其他存储
# glueContext.write_dynamic_frame.from_options(
#     frame=DynamicFrame.fromDF(df, glueContext, "result_df"),
#     connection_type="s3",
#     connection_options={"path": "s3://your-bucket/target-path/"},
#     format="parquet"
# )

job.commit()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 00:45:12