Spark Executor任务分配不均问题排查求助
问题描述
从包含5个分区的Kafka Topic读取数据,因5核无法满足负载需求,将输入重分区至30。Spark进程配置30核,每个Executor分配6核,预期每个Executor会处理6个任务,但实际常出现部分Executor仅处理4个任务、其余处理7个任务的情况,导致作业处理时间倾斜。
作业运行12小时后的Executor指标
| 地址 | 状态 | RDD块 | 存储内存 | 磁盘使用 | 核数 | 活跃任务数 | 失败任务数 | 已完成任务数 | 总任务数 | 任务时间(GC时间) | 输入数据量 | Shuffle读取量 | Shuffle写入量 |
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| ip1:36759 | Active | 7 | 1.6 MB / 144.7 MB | 0.0 B | 6 | 6 | 0 | 442506 | 442512 | 35.9 h (26 min) | 42.1 GB | 25.9 GB | 24.7 GB |
| ip2:36689 | Active | 0 | 0.0 B / 128 MB | 0.0 B | 0 | 0 | 0 | 0 | 0 | 0 ms (0 ms) | 0.0 B | 0.0 B | 0.0 B |
| ip5:44481 | Active | 7 | 1.6 MB / 144.7 MB | 0.0 B | 6 | 6 | 0 | 399948 | 399954 | 29.0 h (20 min) | 37.3 GB | 22.8 GB | 24.7 GB |
| ip1:33187 | Active | 7 | 1.5 MB / 144.7 MB | 0.0 B | 6 | 5 | 0 | 445720 | 445725 | 35.9 h (26 min) | 42.4 GB | 26 GB | 24.7 GB |
| ip3:34935 | Active | 7 | 1.6 MB / 144.7 MB | 0.0 B | 6 | 6 | 0 | 427950 | 427956 | 33.8 h (23 min) | 40.5 GB | 24.8 GB | 24.7 GB |
| ip4:38851 | Active | 7 | 1.7 MB / 144.7 MB | 0.0 B | 6 | 6 | 0 | 410276 | 410282 | 31.6 h (24 min) | 39 GB | 23.9 GB | 24.7 GB |
可见ip5:44481的已完成任务数存在倾斜,且未发现异常GC活动。请问需查看哪些指标来分析该倾斜问题?
更新1
经进一步调试发现,所有分区的数据量不均,但每个任务分配的记录数大致相同。以下是重分区后Stage的统计数据:
| Executor ID | 地址 | 任务时间 | 总任务数 | 失败任务数 | 被杀死任务数 | 成功任务数 | Shuffle读取大小/记录数 | 是否被列入黑名单 |
|---|---|---|---|---|---|---|---|---|
| 0 stdout stderr | ip3:37049 | 0.8 s | 6 | 0 | 0 | 6 | 600.9 KB / 272 | FALSE |
| 1 stdout stderr | ip1:37875 | 0.6 s | 6 | 0 | 0 | 6 | 612.2 KB / 273 | FALSE |
| 2 stdout stderr | ip3:41739 | 0.7 s | 5 | 0 | 0 | 5 | 529.0 KB / 226 | FALSE |
| 3 stdout stderr | ip2:38269 | 0.5 s | 6 | 0 | 0 | 6 | 623.4 KB / 272 | FALSE |
| 4 stdout stderr | ip1:40083 | 0.6 s | 7 | 0 | 0 | 7 | 726.7 KB / 318 | FALSE |
可见任务数与记录数成正比,下一步将排查分区函数的工作机制。
更新2
目前找到的解释是Spark采用轮询分区策略,且该策略在每个分区上独立执行。例如node1有5条记录、node2有7条记录时,node1的轮询会将约3条记录分配给node1、2条给node2;node2的轮询会将约4条记录分配给node1、3条给node2,最终可能出现node1有7条、node2有5条的情况,具体取决于框架代码中每个节点的节点顺序。
注意:性能最优的Executor位于同一IP,是否因为Shuffle后同一主机内的数据传输比跨IP更快?
根据上述数据,重分区已将记录均匀分配至30个分区,但问题是为何部分Executor分配到的分区数更多?
更新3
附上代码片段:
protected JavaDStream<byte[]> getRepartitionedValueStream(JavaInputDStream<ConsumerRecord<String, byte[]>> stream) { return stream.mapToPair(new PairFunction<ConsumerRecord<String, byte[]>, String, byte[]>() { public Tuple2<String, byte[]> call(ConsumerRecord<String, byte[]> x) { return new Tuple2(x.key(), x.value()); } }).groupByKey(this.partitions).flatMap((x) -> { return ((Iterable)x._2()).iterator(); }); }
我们从Kafka获取记录后使用groupByKey,使同一分区的事件归为同一组。
分析与解答
一、需查看的核心指标(针对初始倾斜问题)
- 分区级任务统计:查看30个目标分区各自的执行时间、处理记录数、数据量,确认是否存在单个分区负载远超其他分区的情况,这是任务倾斜的直接来源。
- Shuffle阶段细节:
- Shuffle写入的分区分布:检查每个Task的Shuffle写入量,确认是否因
groupByKey导致热点Key引发分区数据不均。 - Shuffle读取本地性:统计每个Executor读取的本地Shuffle数据占比,验证同一IP下Executor性能更优是否源于本地数据传输优势。
- Shuffle写入的分区分布:检查每个Task的Shuffle写入量,确认是否因
- 任务调度日志:追踪Driver端调度日志,查看是否存在Executor资源波动(如临时CPU抢占)或历史黑名单记录影响任务分配逻辑。
- Kafka源分区数据:检查Kafka Topic的5个源分区的消息堆积量、消费速度,确认是否源分区本身数据不均导致后续倾斜。
二、部分Executor分配更多分区的原因解析
结合groupByKey的逻辑与Spark轮询分区特性,问题根源在于:
- Key哈希分布不均:
groupByKey依赖Key的哈希值分区,若Kafka源数据的Key分布失衡,会导致部分目标分区记录数偏多;Spark调度器会将这些负载高的分区任务分配给可用Executor,进而导致该Executor任务数增加。 - 轮询策略的独立性:Spark的轮询分区在每个源分区独立执行,源分区之间的记录数差异会导致目标分区分布不均,最终反映为Executor任务分配数量不一致。
- 数据本地性调度:Spark会优先将任务分配到数据所在节点的Executor,若多个负载高的分区数据落在同一IP节点,该节点的Executor会被分配更多任务。
三、优化建议
- 替换
groupByKey:若业务允许,改用reduceByKey或aggregateByKey,它们会在Map端做局部聚合,减少Shuffle数据量与热点Key影响。 - 自定义分区器:对热点Key添加随机后缀,将其分散到多个分区,避免单个分区负载过高。
- 调整调度参数:开启
spark.speculation推测执行慢任务;调整spark.locality.wait减少本地性等待时间,均衡任务分配。 - 均衡Kafka源数据:优化生产者分区策略,确保消息均匀分布到Kafka的5个源分区,从源头减少倾斜可能性。
内容的提问来源于stack exchange,提问作者best wishes
相关产品推荐
相关产品推荐

