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

