Spark Streaming分组后仅2个Executor工作其余闲置问题排查
问题根源
核心问题在于分组键的哈希分布与Spark分区数不匹配,结合groupByKey的shuffle特性导致负载不均:
groupByKey通过分组键的哈希值将数据分配到不同分区,每个分区对应一个Task,Task只会被调度到空闲的Executor执行。- 你的场景里有10个唯一unitID,但如果默认分区数过少(或者哈希计算后,这10个键的哈希值仅落到2个分区),就只会生成2个Task,自然只有2个Executor有工作可做,其余闲置。
- 另外,如果Spark Streaming的初始数据源分区数不足(比如Kafka分区数量少),后续分组操作的分区数会继承这个限制,进一步加剧负载不均。
实现任务均衡分配的解决方法
1. 显式指定分组后的分区数
直接在groupByKey后添加repartition,将分区数设置为Executor数量(或略大于),让10个unitID分散到更多分区:
// 分组后重分区到8个,与Executor数量对齐 val dfGrouped: KeyValueGroupedDataset[Int, Car] = dfSource.groupByKey(car1 => car1.unitID).repartition(8)
这样每个分区会分配1-2个unitID的任务,所有Executor都能被调度到对应Task。
2. 全局配置Shuffle分区数
在SparkSession初始化时设置spark.sql.shuffle.partitions,让所有需要shuffle的操作(包括groupByKey)默认生成足够的分区:
val spark = SparkSession.builder() .appName("CarStreamingApp") .config("spark.sql.shuffle.partitions", "8") // 与Executor数量匹配 .getOrCreate()
这个配置可以避免后续所有shuffle操作出现分区数不足的问题。
3. 提前对数据源分区
如果初始数据源的分区数过少,先对输入DataFrame重分区,再执行分组操作:
// 先将数据源分区到8个,再执行分组逻辑 val dfSourceRepartitioned = dfSource.repartition(8) val dfGrouped = dfSourceRepartitioned.groupByKey(car1 => car1.unitID)
内容的提问来源于stack exchange,提问作者MAK
相关产品推荐
相关产品推荐

