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

如何在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状态码等)。

具体步骤

  1. 优先捕获Py4JJavaError异常,而非泛泛的Exception
  2. 从异常对象中获取Java端抛出的底层异常实例
  3. 判断是否为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 19:25:21