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

AWS Glue批量抽取Oracle表时连接耗尽问题及优化咨询

问题场景

使用AWS Glue通过循环从Oracle数据库抽取约500张表到S3,前200多张表抽取成功,后续出现如下报错:

An error occurred while calling o41372.getDynamicFrame. : java.sql.SQLRecoverableException: IO Error: Got minus one from a read call, connect lapse 22 ms., Authentication lapse 0 ms.
at oracle.jdbc.driver.T4CConnection.handleLogonIOException(T4CConnection.java:935)
at oracle.jdbc.driver.T4CConnection.logon(T4CConnection.java:700)
at oracle.jdbc.driver.PhysicalConnection.connect(PhysicalConnection.java:1041)
at oracle.jdbc.driver.T4CDriverExtension.getConnection(T4CDriverExtension.java:89)
at oracle.jdbc.driver.OracleDriver.connect(OracleDriver.java:732)
at oracle.jdbc.driver.OracleDriver.connect(OracleDriver.java:648)
at com.amazonaws.services.glue.util.JDBCWrapper$.$anonfun$connectionProperties$5(JDBCUtils.scala:1120)
at com.amazonaws.services.glue.util.JDBCWrapper$.$anonfun$connectWithSSLAttempt$2(JDBCUtils.scala:1071)
at scala.Option.getOrElse(Option.scala:189)
at com.amazonaws.services.glue.util.JDBCWrapper$.$anonfun$connectWithSSLAttempt$1(JDBCUtils.scala:1071)
at scala.Option.getOrElse(Option.scala:189)
at com.amazonaws.services.glue.util.JDBCWrapper$.connectWithSSLAttempt(JDBCUtils.scala:1071)
at com.amazonaws.services.glue.util.JDBCWrapper$.connectionProperties(JDBCUtils.scala:1116)
at com.amazonaws.services.glue.util.JDBCWrapper.connectionProperties$lzycompute(JDBCUtils.scala:822)

推测报错原因是达到Oracle最大连接数限制,Oracle强制关闭后续连接。现有代码片段如下:

def initial_extract_data_to_s3(CONNECTION,DATABASE,SCHEMA,TABLE):
    try:
        data_df = glueContext.create_dynamic_frame.from_options(
                                    connection_type = "oracle",
                                    connection_options = {
                                        "useConnectionProperties": "true",
                                        "dbtable": f"{schema_name}.{table_name}",
                                        "sampleQuery":src_query,
                                        "connectionName": f"{connection}",
                                        "hashexpression":f"{partition_key}",
                                        "hashpartitions":"12",
                                    },
                                    transformation_ctx = "data_df")

        
        data_df.repartition(1).write.mode("overwrite").option("compression", "snappy").parquet(f"s3://{db_name}/{schema_name}/{table_name}")
            
    except Exception as exc: 
        print(exc)

# Calling the above function below:

for row in df.collect():
    try:
        connection = row['CONNECTION']
        db_name = row['DATABASE']
        schema_name = row['SCHEMA']
        table_name = row['TABLE']

        # Call your functions here
        initial_extract_data_to_s3(CONNECTION,DATABASE,SCHEMA,TABLE)

    except Exception as e:
        continue  # Continue to the next table

需求:能否在AWS Glue中打开新连接前主动关闭旧连接?除了将作业拆分为每200张表一组,有没有更好的方法在单个作业内完成500张表的抽取?


解决方案

1. 显式释放Spark资源,加速连接回收

AWS Glue的DynamicFrame底层依赖Spark DataFrame,默认情况下Spark会缓存DataFrame和关联的JDBC连接。可以在每张表处理完成后,手动释放资源:

def initial_extract_data_to_s3(CONNECTION,DATABASE,SCHEMA,TABLE):
    try:
        data_df = glueContext.create_dynamic_frame.from_options(
                                    connection_type = "oracle",
                                    connection_options = {
                                        "useConnectionProperties": "true",
                                        "dbtable": f"{schema_name}.{table_name}",
                                        "sampleQuery":src_query,
                                        "connectionName": f"{connection}",
                                        "hashexpression":f"{partition_key}",
                                        "hashpartitions":"12",
                                    },
                                    transformation_ctx = "data_df")

        data_df.repartition(1).write.mode("overwrite").option("compression", "snappy").parquet(f"s3://{db_name}/{schema_name}/{table_name}")
        
        # 显式释放资源,加速连接回收
        data_df.toDF().unpersist()  # 释放DataFrame内存及关联的JDBC连接
        spark.catalog.clearCache()  # 清理Spark全局缓存
            
    except Exception as exc: 
        print(exc)

2. 手动管理JDBC连接,确保用完即关

绕过Glue的自动连接管理,手动创建和关闭Oracle连接,彻底控制连接生命周期:

from py4j.java_gateway import java_import
java_import(spark._jvm, "oracle.jdbc.driver.OracleDriver")

def get_oracle_connection(connection_name):
    # 从Glue连接中提取JDBC配置信息
    glue_conn_conf = glueContext.extract_jdbc_conf(connection_name)
    db_url = glue_conn_conf["url"]
    db_user = glue_conn_conf["user"]
    db_pwd = glue_conn_conf["password"]
    # 创建Oracle JDBC连接
    conn = spark._jvm.java.sql.DriverManager.getConnection(db_url, db_user, db_pwd)
    return conn

def initial_extract_data_to_s3(conn, db_name, schema_name, table_name, partition_key):
    try:
        # 通过自定义连接读取数据
        query = f"SELECT * FROM {schema_name}.{table_name}"
        spark_df = spark.read.jdbc(
            url=conn.getMetaData().getURL(),
            table=f"({query}) AS temp_table",
            properties={
                "user": conn.getMetaData().getUserName(),
                "password": conn.getMetaData().getPassword()
            }
        )
        # 写入S3
        spark_df.repartition(1).write.mode("overwrite").option("compression", "snappy").parquet(f"s3://{db_name}/{schema_name}/{table_name}")
    except Exception as exc:
        print(f"抽取表{schema_name}.{table_name}失败: {exc}")
    finally:
        # 强制关闭连接,避免连接堆积
        if conn:
            conn.close()

# 循环调用逻辑修改
for row in df.collect():
    try:
        connection = row['CONNECTION']
        db_name = row['DATABASE']
        schema_name = row['SCHEMA']
        table_name = row['TABLE']
        partition_key = row.get('PARTITION_KEY', '')  # 适配你的表分区键字段

        # 创建连接并处理表
        oracle_conn = get_oracle_connection(connection)
        initial_extract_data_to_s3(oracle_conn, db_name, schema_name, table_name, partition_key)
    except Exception as e:
        print(f"处理表{table_name}出错: {e}")
        continue

3. 优化Glue连接参数,减少连接占用

在connection_options中添加以下参数,优化连接行为:

  • batchSize: 设置批量读取行数,减少单连接的占用时间,例如"batchSize": "20000"
  • connectionProperties: 添加Oracle JDBC参数,控制连接超时和回收,例如:
    "connectionProperties": "oracle.net.CONNECT_TIMEOUT=15000;oracle.net.READ_TIMEOUT=60000;oracle.jdbc.ReadTimeout=60000"
    
  • 若不是必须,暂时移除hashpartitions配置,减少并发连接数

4. 调整Oracle端配置(如果有权限)

  • 检查并调大Oracle的processes和sessions参数,增加允许的最大连接数
  • 配置Oracle连接超时机制,让闲置连接自动回收,例如设置IDLE_TIME参数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 08:12:09