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
相关产品推荐
相关产品推荐

