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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 12:45:36