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

Spark读取Oracle大表写入文件超时问题求助

问题描述

环境与需求

  • Spark版本:2.4.7
  • 操作系统:Linux RHEL/Fedora
  • 核心需求:每日读取Oracle中约1.5亿条记录的表,重分区为800份后写入共享文件系统

异常现象

  • 小数据量测试(100万条、重分区6)时,读写仅需2-3分钟,完全正常
  • 全量运行时,读取DataFrame没问题,但写入操作始终卡着完不成
  • 日志反复刷这条:INFO sort.UnsafeExternalSorter: Thread 135 spilling sort data of 5.2 GB to disk (57 times so far)
  • 跑了2.5小时后,Shuffle Spill统计:内存259.4GB,磁盘22.0GB

提交命令与核心代码

提交命令

time spark-submit --verbose --conf spark.dynamicAllocation.enabled=false --conf spark.spark.sql.broadcastTimeout=1000 --conf spark.sql.shuffle.partitions=1500 --conf "spark.ui.enabled=true" --master yarn --driver-memory 60G --num-executors 40 --executor-memory 10G --executor-cores 8 --jars ojdbc6.jar SparkOracleExtractor.jar

核心代码

val myquery = "select * from mytable"
val dt=20221023
spark.read.format("jdbc")
  .option("url", s"jdbc:oracle:thin:@$db_connect_string")
  .option("driver", "oracle.jdbc.driver.OracleDriver")
  .option("query", myquery)
  .option("user", db_user)
  .option("password", db_pass)
  .option("fetchsize", 10000)
  .option("delimiter", "|")
  .load()
  .repartition(800)
  .write.csv(s"file:///fs/extrat_path/${dt}")

集群状态

集群资源充足,队列使用率仅5.0%,必须用Spark解决问题并提速


解决方案与优化建议

1. 从JDBC读取阶段就实现并行,彻底避免全量Shuffle

你现在的问题根源是:JDBC默认单分区拉取1.5亿条数据,之后执行repartition(800)会触发全量Shuffle——单分区的数据要打散到800个分区,必然导致大量内存溢出刷盘,这就是任务卡主的核心原因。

直接在JDBC读取时就拆成800个并行分区,后续不需要再做repartition:

// 假设表有数值型的主键/分区字段(比如id),用它来拆分查询
val df = spark.read.format("jdbc")
  .option("url", s"jdbc:oracle:thin:@$db_connect_string")
  .option("driver", "oracle.jdbc.driver.OracleDriver")
  .option("dbtable", "mytable")
  .option("user", db_user)
  .option("password", db_pass)
  .option("fetchsize", 10000) // 保持这个参数,提升单分区读取效率
  .option("partitionColumn", "id") // 选一个均匀分布的数值/日期字段
  .option("lowerBound", "1") // 该字段的最小值
  .option("upperBound", "150000000") // 该字段的最大值(对应1.5亿条数据)
  .option("numPartitions", 800) // 直接拆成目标分区数
  .load()

如果没有合适的数值字段,也可以手动拆分SQL,用query参数配合分段条件实现并行(比如按日期分段、或按id范围手动拆分800段)。

2. 调整Shuffle相关配置(如果必须保留Shuffle操作)

如果因为某些原因没法从读取阶段并行,就调整以下配置缓解Spill压力:

  • 增大Shuffle可用内存占比:--conf spark.shuffle.memoryFraction=0.4(默认0.2,Spark 2.4中可调整,让Shuffle时能用到更多Executor内存)
  • 确保Spill到磁盘的数据是压缩的:--conf spark.shuffle.spill.compress=true(默认开启,但确认下没被修改)
  • 对齐Shuffle分区和目标写入分区:把spark.sql.shuffle.partitions从1500改成800,减少不必要的中间分区数

3. 优化Executor资源配置

当前--executor-cores 8,每个Executor扛的任务数太多,容易出现内存竞争。在集群资源充足的前提下,调整为:

--num-executors 80 --executor-cores 4 --executor-memory 10G

把Executor数量翻倍,每个Executor的core数减半,提升整体并行度,减少单个Executor的压力。

4. 写入阶段的小优化

  • 控制单个CSV文件的大小:添加option("maxRecordsPerFile", "187500")(1.5亿/800=187500),避免单个文件过大,提升写入稳定性
  • 改用HDFS路径替代file://:如果是YARN集群,不要用本地文件系统路径,改成HDFS路径(比如hdfs:///fs/extrat_path/${dt}),避免每个Executor写入本地后再同步的瓶颈

5. 其他细节优化

  • 升级JDBC驱动:用ojdbc8.jar替换ojdbc6.jar,Oracle 12c+对ojdbc8的支持更好,性能更优
  • 不要读全表字段:如果不需要select *,改成具体需要的字段列表,减少数据传输量和内存占用

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 06:15:37