Spark写入S3时无法充分利用Worker节点的问题排查
Spark从Teradata导出数据到S3时核心利用率低的问题分析与优化
Spark核心数量的决策逻辑
Spark的核心(Executor核心)利用率完全由活跃任务数量决定,而任务数量直接对应RDD/DataFrame的分区数:
- 初始阶段用满12核,是因为JDBC读取Teradata的分区任务在并行执行,每个分区对应一个任务,占满所有核心
- 后续只剩3核,说明写入S3阶段的任务并行度被限制,活跃任务数降到3个,导致大量核心闲置
- 核心分配的底层逻辑:每个Worker节点的Executor数量由
spark.executor.cores和节点资源决定,你的集群3个Worker×4核,建议设置spark.executor.cores=4,每个Worker启动1个Executor,总核心数12,确保资源分配最大化
写入S3并行度不足的优化方案(除S3A Committer外)
1. 确保写入阶段的分区数匹配核心数
- 先排查读取后的数据分区数:在写入S3前执行
println(df.rdd.getNumPartitions()),确认分区数是否保持12或更高- 如果分区数不足,说明JDBC读取的分区配置没生效,或者中间操作(如
coalesce、聚合)导致分区减少 - 显式用
repartition(n)强制调整分区数,35亿数据建议设置为24-36个分区(略大于总核心数,应对IO瓶颈),注意repartition会触发shuffle,需根据数据量合理设置
- 如果分区数不足,说明JDBC读取的分区配置没生效,或者中间操作(如
- 检查JDBC读取的分区策略:Teradata的JDBC分区是否均匀?如果某个分区数据量远大于其他,会导致后续写入阶段该任务占用大量资源,其他核心闲置(数据倾斜)
- 解决数据倾斜:对分区字段加盐拆分任务,或调整JDBC的分区列(选择分布均匀的字段),确保每个分区数据量均衡
2. 排查任务调度与资源瓶颈
- 打开Spark UI的Stages页面,查看写入阶段的任务数、执行时间、任务状态:
- 如果只有少数任务在跑,优先排查数据倾斜:查看每个任务的输入数据量,找到倾斜的分区
- 查看Executors页面,检查是否有Executor被kill、GC频繁:
- 调整
spark.executor.memoryOverhead,建议设置为Executor总内存的10%-20%(你的场景每个Executor内存256GB,可设为20-30GB),避免内存不足导致Executor退出
- 调整
- 启用自适应执行:设置
spark.sql.adaptive.enabled=true,让Spark根据实际数据分布自动调整分区数和并行度,应对动态变化的数据量
3. 优化S3写入的IO配置
- 调整
spark.sql.shuffle.partitions:如果写入前有shuffle操作,该参数控制shuffle后的分区数,建议设置为与总核心数匹配的值(如12/24) - 增加S3连接数:设置
spark.hadoop.fs.s3a.connection.maximum=100(默认是15),避免IO连接不足导致任务阻塞 - 启用S3快速上传:设置
spark.hadoop.fs.s3a.fast.upload=true,启用分块并行上传,提升写入吞吐量 - 确认Parquet压缩配置:确保
spark.sql.parquet.compression.codec=snappy生效,减少写入的数据量,降低IO压力
4. 循环导出场景的特殊优化(场景2)
- 每次循环导出后,显式清理数据缓存:执行
df.unpersist(),避免内存累积导致Executor资源不足 - 每次循环重新创建JDBC读取的DataFrame,避免复用旧的RDD/DataFrame导致分区数意外减少
- 循环内的任务可以设置
spark.sql.shuffle.partitions为与单次任务匹配的数值,避免资源浪费
历史案例的补充分析
之前3个1核10GB Worker处理5亿数据时,先2个Worker后1个的情况,同样是并行度下降导致:
- 要么是读取的分区数不足,要么是写入阶段数据倾斜或分区被合并,导致活跃任务数减少
- 优化思路和上述一致,优先排查分区数和数据分布
内容的提问来源于stack exchange,提问作者the_data_novice
相关产品推荐
相关产品推荐

