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
相关产品推荐
相关产品推荐

