Spark作业添加distinct后频繁失败、运行超时的原因咨询
问题场景
原有Spark作业流程:从Delta Table-1读取数据,执行基础Group By聚合后写入Delta Table-2,运行稳定。后续发现Table-1存在约20%重复数据,遂在读取阶段添加Distinct操作去重,SQL代码如下:
create or replace temporary view unique_data as select distinct a.*, <few minor transforms> from <delta table-1> a;
添加后作业频繁因内存问题失败,运行时长增至原作业的5-6倍,报错信息如下:
Py4JJavaError: An error occurred while calling o421.saveAsTable. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 15 in stage 47.0 failed 4 times, most recent failure: Lost task 15.3 in stage 47.0 (TID 1850) (10.35.5.3 executor 22): ExecutorLostFailure (executor 22 exited caused by one of the running tasks) Reason: Command exited with code 52 Driver stacktrace: at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:3950) at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:3868) at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:3855) at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62) at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55) at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49) at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:3855) at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1716) at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1699) at scala.Option.foreach(Option.scala:407) at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1699) at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:4196) at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:4108) at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:4096) at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:54)
切换至n2-highmem-16实例(每Worker 128GB内存,共5个Worker,开启自动扩缩容)后,作业恢复正常。
核心问题解答
Distinct属于高开销操作吗?
是的,Distinct本质上属于高开销操作,它和Group By的底层执行逻辑高度相似——都是基于哈希分区做Shuffle,再在Executor端通过哈希表去重/聚合,但你的场景中几个关键因素放大了开销,直接导致作业失败:
Shuffle数据量暴增
原作业的Group By仅基于聚合键做Shuffle,数据量小;而你添加的Distinct是基于a.*+转换后的所有字段做哈希分区,需要把Table-1的全量数据(含20%重复)都进行Shuffle,数据量远大于原作业,直接拉高了网络IO和磁盘IO开销。Executor内存压力过载
Distinct的去重逻辑依赖Executor内存中的哈希表存储已处理过的完整记录。原作业的Group By仅需存储聚合键的中间结果,键的数量远小于全量记录数;而Distinct需要存储去重后的80%全量记录,内存占用量呈数量级增长,直接撑爆原有Executor内存,导致进程崩溃(报错中的ExecutorLostFailure就是内存不足引发的进程退出)。数据倾斜与重试放大耗时
全字段哈希分区更容易引发数据倾斜——部分分区的重复记录集中,对应Executor需要处理远超平均量的数据,内存不足导致任务反复失败重试,进一步拉长了整体运行时间,最终变成原作业的5-6倍。
优化建议
- 缩小Distinct字段范围:如果存在唯一标识字段(如主键),仅对唯一键做Distinct,再关联原表获取其他字段,大幅减少Shuffle和内存占用。
- 调整内存参数:增大
spark.executor.memory和spark.executor.memoryOverhead,给Executor预留更多内存存储去重哈希表。 - 优化Shuffle配置:调整
spark.sql.shuffle.partitions增加分区数,缓解数据倾斜;开启spark.sql.adaptive.enabled让Spark自动适配资源与分区。
内容的提问来源于stack exchange,提问作者steve

