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

AWS Glue中Spark JDBC读取Oracle表过慢问题求助

优化Spark JDBC读取Oracle大表并关联查询的方案

看起来你遇到的核心问题是Spark JDBC的分区策略不合理导致数据倾斜,加上不必要的全表读取,使得数据加载远慢于Oracle原生执行。结合你的场景,我整理了几个针对性的优化方案:

1. 优先让Oracle完成关联聚合(最快的解决方案)

既然SQL Developer执行关联查询仅需25分钟,说明Oracle的优化器已经能高效处理这个查询。与其把5亿条数据拉到Spark再做关联,直接让Oracle执行完整查询,Spark只读取最终结果,这会大幅减少数据传输量和Spark的计算压力。

修改你的JDBC读取代码如下:

val oracleQuery = """
SELECT CP.CLIENT_ID, COUNT(1) NoofCases 
FROM FSP.CUSTOMER_CASE CC 
JOIN FSP.GROUP FG ON FG.ID = CC.OWNER_ID 
JOIN FSP.CLIENT_PLATFORM CP ON CP.CLIENT_ID = SUBSTR(FG.PATH, 2, INSTR(FG.PATH, '/') + INSTR(SUBSTR(FG.PATH, 1 + INSTR(FG.PATH, '/')), '/') - 2) 
WHERE FG.STATUS = 'ACTIVE' AND FG.TYPE = 'CLIENT' 
GROUP BY CP.CLIENT_ID
"""

val resultDf = glueContext.read.format("jdbc") 
  .option("url","jdbc:oracle:thin://abcd:1521/abcd.com") 
  .option("user","USER_PROD") 
  .option("password","ffg#Prod") 
  .option("query", oracleQuery) 
  .option("driver","oracle.jdbc.OracleDriver")
  .option("fetchSize", "10000") // 增加fetchSize减少网络往返次数
  .load()

这个方案的性能会和SQL Developer接近,因为所有复杂计算都在Oracle端完成,Spark只需要拉取最终的聚合结果。

2. 优化Spark JDBC读取的分区策略

如果必须要把数据拉到Spark处理,那首先要解决分区列数据倾斜的问题:
你选择的OUTSTANDING_ACTIONS列数据分布极不均匀——值0和1的记录数远大于其他值,这会导致对应分区的任务处理量远超其他分区,大部分Executor闲置,整体拖慢加载速度。

优化措施:

  • 更换均匀分布的分区列:优先选择自增主键、时间戳(如果有按时间均匀分布的字段),或者Oracle的ROWID(可转换为数值进行范围分区)。例如,如果表有自增ID列CASE_ID:
df = glueContext.read.format("jdbc") 
  .option("url","jdbc:oracle:thin://abcd:1521/abcd.com") 
  .option("user","USER_PROD") 
  .option("password","ffg#Prod") 
  .option("numPartitions", 80) // 根据worker数量调整,40个worker可设为80-120
  .option("partitionColumn", "CASE_ID") 
  .option("lowerBound", 1) 
  .option("upperBound", 500000000) // 对应5亿条记录的ID上限
  .option("fetchSize", "10000")
  .option("dbtable","FSP.CUSTOMER_CASE") 
  .option("driver","oracle.jdbc.OracleDriver").load()
  • 如果没有均匀列,用自定义分区:比如通过MOD(CASE_ID, numPartitions)来拆分数据,确保每个分区数据量接近:
// 手动生成多个查询,每个查询对应一个分区
val numPartitions = 80
val partitionQueries = (0 until numPartitions).map(i => 
  s"(SELECT * FROM FSP.CUSTOMER_CASE WHERE MOD(CASE_ID, $numPartitions) = $i) AS part_$i"
).mkString(",")

df = glueContext.read.format("jdbc") 
  .option("url","jdbc:oracle:thin://abcd:1521/abcd.com") 
  .option("user","USER_PROD") 
  .option("password","ffg#Prod") 
  .option("dbtable", partitionQueries)
  .option("fetchSize", "10000")
  .option("driver","oracle.jdbc.OracleDriver").load()

3. 减少不必要的数据读取

你的关联查询只用到了CUSTOMER_CASE表的OWNER_ID列,没必要读取全表所有字段。修改dbtable或query只读取需要的列,能大幅减少数据传输量:

df = glueContext.read.format("jdbc") 
  .option("url","jdbc:oracle:thin://abcd:1521/abcd.com") 
  .option("user","USER_PROD") 
  .option("password","ffg#Prod") 
  .option("numPartitions", 80)
  .option("partitionColumn", "CASE_ID") 
  .option("lowerBound", 1) 
  .option("upperBound", 500000000)
  .option("fetchSize", "10000")
  .option("dbtable","(SELECT OWNER_ID FROM FSP.CUSTOMER_CASE) AS CC_SUB") // 仅读取需要的列
  .option("driver","oracle.jdbc.OracleDriver").load()

4. Glue环境的额外优化

  • 调整Executor配置:虽然你换成了最大的GP2型,但要确保每个Executor的内存和CPU匹配数据处理需求。比如r5.4xlarge(16vCPU, 64GB内存)的话,可以设置--executor-memory 40G --executor-cores 8,避免GC频繁导致的性能损耗。
  • 检查网络延迟:如果Oracle是本地部署,确保Glue和Oracle之间的网络连接稳定,考虑用AWS Direct Connect或VPC peering减少跨网络的延迟。
  • 升级Glue版本:使用最新的Glue版本,通常会包含Spark和JDBC驱动的性能优化。

内容的提问来源于stack exchange,提问作者Atharv Thakur

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 18:27:46