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

Spark Executor任务分配不均问题排查求助

问题描述

从包含5个分区的Kafka Topic读取数据,因5核无法满足负载需求,将输入重分区至30。Spark进程配置30核,每个Executor分配6核,预期每个Executor会处理6个任务,但实际常出现部分Executor仅处理4个任务、其余处理7个任务的情况,导致作业处理时间倾斜。

作业运行12小时后的Executor指标

地址状态RDD块存储内存磁盘使用核数活跃任务数失败任务数已完成任务数总任务数任务时间(GC时间)输入数据量Shuffle读取量Shuffle写入量
ip1:36759Active71.6 MB / 144.7 MB0.0 B66044250644251235.9 h (26 min)42.1 GB25.9 GB24.7 GB
ip2:36689Active00.0 B / 128 MB0.0 B000000 ms (0 ms)0.0 B0.0 B0.0 B
ip5:44481Active71.6 MB / 144.7 MB0.0 B66039994839995429.0 h (20 min)37.3 GB22.8 GB24.7 GB
ip1:33187Active71.5 MB / 144.7 MB0.0 B65044572044572535.9 h (26 min)42.4 GB26 GB24.7 GB
ip3:34935Active71.6 MB / 144.7 MB0.0 B66042795042795633.8 h (23 min)40.5 GB24.8 GB24.7 GB
ip4:38851Active71.7 MB / 144.7 MB0.0 B66041027641028231.6 h (24 min)39 GB23.9 GB24.7 GB

可见ip5:44481的已完成任务数存在倾斜,且未发现异常GC活动。请问需查看哪些指标来分析该倾斜问题?


更新1

经进一步调试发现,所有分区的数据量不均,但每个任务分配的记录数大致相同。以下是重分区后Stage的统计数据:

Executor ID地址任务时间总任务数失败任务数被杀死任务数成功任务数Shuffle读取大小/记录数是否被列入黑名单
0
stdout
stderr
ip3:370490.8 s6006600.9 KB / 272FALSE
1
stdout
stderr
ip1:378750.6 s6006612.2 KB / 273FALSE
2
stdout
stderr
ip3:417390.7 s5005529.0 KB / 226FALSE
3
stdout
stderr
ip2:382690.5 s6006623.4 KB / 272FALSE
4
stdout
stderr
ip1:400830.6 s7007726.7 KB / 318FALSE

可见任务数与记录数成正比,下一步将排查分区函数的工作机制。


更新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性能更优是否源于本地数据传输优势。
  • 任务调度日志:追踪Driver端调度日志,查看是否存在Executor资源波动(如临时CPU抢占)或历史黑名单记录影响任务分配逻辑。
  • Kafka源分区数据:检查Kafka Topic的5个源分区的消息堆积量、消费速度,确认是否源分区本身数据不均导致后续倾斜。

二、部分Executor分配更多分区的原因解析

结合groupByKey的逻辑与Spark轮询分区特性,问题根源在于:

  1. Key哈希分布不均:groupByKey依赖Key的哈希值分区,若Kafka源数据的Key分布失衡,会导致部分目标分区记录数偏多;Spark调度器会将这些负载高的分区任务分配给可用Executor,进而导致该Executor任务数增加。
  2. 轮询策略的独立性:Spark的轮询分区在每个源分区独立执行,源分区之间的记录数差异会导致目标分区分布不均,最终反映为Executor任务分配数量不一致。
  3. 数据本地性调度:Spark会优先将任务分配到数据所在节点的Executor,若多个负载高的分区数据落在同一IP节点,该节点的Executor会被分配更多任务。

三、优化建议

  • 替换groupByKey:若业务允许,改用reduceByKey或aggregateByKey,它们会在Map端做局部聚合,减少Shuffle数据量与热点Key影响。
  • 自定义分区器:对热点Key添加随机后缀,将其分散到多个分区,避免单个分区负载过高。
  • 调整调度参数:开启spark.speculation推测执行慢任务;调整spark.locality.wait减少本地性等待时间,均衡任务分配。
  • 均衡Kafka源数据:优化生产者分区策略,确保消息均匀分布到Kafka的5个源分区,从源头减少倾斜可能性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 19:09:21