如何在AWS Glue Python作业中捕获执行存储过程的java.sql.SQLException
问题描述
运行AWS Glue作业执行Oracle数据库中的存储过程时,希望在存储过程失败时捕获数据库特定的SQL异常。当前使用py4j建立连接执行SQL命令,但通过Python的try-except捕获到的只是完整的py4j.protocol.Py4JJavaError错误信息,无法直接提取到Oracle返回的具体错误(比如错误码、原生错误消息)。
想确认是否可以通过java_import(sc._gateway.jvm,"java.sql.SQLException")来从execute()调用中提取这些数据库特定错误,当前的代码如下:
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 import boto3 ############################################ ### UPDATE THE STORE PROCEDURE NAME HERE ### sp_names = [ 'DEV_WF_POC1', 'DEV_WF_POC2' ] ############################################ #Set the conection name (Will be replaced for FT and PROD by powershell deployment script) glue_connection_name = 'dw-dev-connection' #Use systems args to return job name and pass to local variable args = getResolvedOptions(sys.argv, ['JOB_NAME','WORKFLOW_NAME', 'WORKFLOW_RUN_ID']) workflow_name = args['WORKFLOW_NAME'] workflow_run_id = args['WORKFLOW_RUN_ID'] glue_job_name = args['JOB_NAME'] #Create spark handler and update status of glue job sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(glue_job_name, args) job.commit() logger = glueContext.get_logger() glue_client = boto3.client('glue') #Extract connection details from Data Catelog source_jdbc_conf = glueContext.extract_jdbc_conf(glue_connection_name) #Import Python for Java Java.sql libs 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") #Extract the URL from the JDBC connection oracleurl = source_jdbc_conf.get('url') # Update connection string to expected Oracle format oracleurl = oracleurl.replace("oracle://", "oracle:thin:@") oracleurl = oracleurl + ':orcl' #Create the connection to the Oracle database with java.sql conn = sc._gateway.jvm.DriverManager.getConnection(oracleurl, source_jdbc_conf.get('user'), source_jdbc_conf.get('password')) #Change autocommit to false to avoid Table lock error conn.setAutoCommit(False); # error dict errs = {} err = '' try: for sp_name in sp_names: #Prepare call stored procedure statement and execute cstmt = conn.prepareCall("{call reporting." + sp_name + "}"); results = cstmt.execute(); conn.commit(); # capture error except Exception as e: # work on python 3.x ##errs['msg'] = str(sc._gateway.jvm.SQLException.getMessage())- doesn't work errs['error'] = str(e) errs['sp_name'] = sp_name errs['error_type'] = str(type(e)).replace("<class '","").replace("'>","") if len(errs) != 0: stmt = conn.createStatement(); sql = "insert into dev_workflow_errors (timestamp, workflow_id, workflow_name, job_name, sp_name, error_type, error) values (current_timestamp, '" + workflow_run_id + "', '" + workflow_name + "', '" + glue_job_name + "', '" + errs['sp_name'] + "', '" + errs['error_type'] + "', '" + errs.get('msg','') + "')" rs = stmt.executeUpdate(sql); conn.commit(); #sys.exit(1) #Close down the connection conn.close(); #Update Logger logger.info("Finished")
解决方案
可以通过py4j提供的机制,从Py4JJavaError中解析出底层的SQLException,进而提取Oracle返回的数据库特定错误信息(比如错误码、原生错误消息、SQL状态码等)。
具体步骤
- 优先捕获
Py4JJavaError异常,而非泛泛的Exception - 从异常对象中获取Java端抛出的底层异常实例
- 判断是否为
SQLException类型,调用其Java方法提取具体错误信息
修改后的异常捕获代码
替换原有except块为以下内容:
from py4j.protocol import Py4JJavaError # ... 其余代码保持不变 try: for sp_name in sp_names: cstmt = conn.prepareCall("{call reporting." + sp_name + "}") results = cstmt.execute() conn.commit() except Py4JJavaError as e: # 获取Java端的原始异常对象 java_exception = e.java_exception # 判断是否为SQL异常 if isinstance(java_exception, sc._gateway.jvm.java.sql.SQLException): errs['sp_name'] = sp_name errs['error_type'] = 'SQLException' errs['error_msg'] = java_exception.getMessage() errs['error_code'] = java_exception.getErrorCode() # Oracle专属错误码(如ORA-xxxx) errs['sql_state'] = java_exception.getSQLState() # SQL标准状态码 else: # 处理其他Java异常 errs['sp_name'] = sp_name errs['error_type'] = str(type(java_exception)).replace("class ","") errs['error_msg'] = str(java_exception) except Exception as e: # 处理非Py4J的Python本地异常 errs['sp_name'] = sp_name errs['error_type'] = str(type(e)).replace("<class '","").replace("'>","") errs['error_msg'] = str(e)
错误日志插入优化
为避免错误消息中的特殊字符(如单引号)导致SQL语法错误,建议使用参数化查询插入日志:
if len(errs) != 0: stmt = conn.prepareStatement(""" insert into dev_workflow_errors ( timestamp, workflow_id, workflow_name, job_name, sp_name, error_type, error, error_code, sql_state ) values ( current_timestamp, ?, ?, ?, ?, ?, ?, ?, ? ) """) stmt.setString(1, workflow_run_id) stmt.setString(2, workflow_name) stmt.setString(3, glue_job_name) stmt.setString(4, errs['sp_name']) stmt.setString(5, errs['error_type']) stmt.setString(6, errs.get('error_msg', '')) stmt.setInt(7, errs.get('error_code', 0)) stmt.setString(8, errs.get('sql_state', '')) stmt.executeUpdate() conn.commit() # 根据业务需求决定是否终止作业 sys.exit(1)
关键说明
e.java_exception:直接获取Py4J包装的Java原生异常对象isinstance(...):精准判断异常类型,确保只处理数据库相关的SQL异常getErrorCode():返回Oracle专属错误码(如ORA-00001),是定位数据库问题的核心信息getMessage():返回Oracle原生的错误描述文本,包含具体的失败原因
内容的提问来源于stack exchange,提问作者Marcus Trinder
相关产品推荐
相关产品推荐

