在 Spark 中,**分区器(Partitioner)**决定了 RDD 中每条 Key 会被分配到哪一个分区,主要影响 groupByKeyreduceByKeyjoinShuffle 类算子 的行为。

Spark 内置的分区器主要有 三种,另外还支持 自定义分区器


一、Spark 内置的三种分区器

1️⃣ HashPartitioner(默认)

类: org.apache.spark.HashPartitioner

核心原理
partitionId = key.hashCode() % numPartitions
特点
  • 最常用、默认分区器
  • 相同 key 一定进入同一个分区
  • 实现简单、开销小
  • 可能导致数据倾斜(hash 分布不均)
使用场景
  • reduceByKeygroupByKeydistinct
  • 不关心 key 的分布顺序
  • 大多数通用场景
示例
val rdd = sc.parallelize(Seq(("a",1), ("b",2), ("c",3)))
val partitioned = rdd.partitionBy(new HashPartitioner(3))

2️⃣ RangePartitioner(范围分区器)

类: org.apache.spark.RangePartitioner

核心原理
  1. 先对 key 进行采样
  2. 按 key 的大小排序
  3. 将 key 范围划分成 numPartitions 个区间
  4. 每个区间一个分区
特点
  • key 是有序的
  • 尽量保证 每个分区数据量均匀
  • 会产生 额外采样开销
  • 适合排序、范围查询
使用场景
  • sortByKey
  • repartitionAndSortWithinPartitions
  • 需要全局有序或范围查询
示例
val rdd = sc.parallelize(Seq((1,"a"), (3,"b"), (2,"c")))
val partitioned = rdd.partitionBy(new RangePartitioner(3, rdd))

3️⃣ CustomPartitioner(自定义分区器)

类: 实现 org.apache.spark.Partitioner

核心原理
  • 用户自己定义 key → partitionId 的映射规则
特点
  • 完全可控
  • 可解决 数据倾斜
  • 可结合业务规则(如按城市、按用户 ID 段)
使用场景
  • 数据倾斜严重
  • 业务规则强相关
  • 特殊 key 需要单独处理
示例
class MyPartitioner(numParts: Int) extends Partitioner {
  override def numPartitions: Int = numParts

  override def getPartition(key: Any): Int = {
    key match {
      case k: String if k.startsWith("hot") => 0
      case _ => math.abs(key.hashCode % numPartitions)
    }
  }
}

rdd.partitionBy(new MyPartitioner(4))

二、三种分区器对比总结

分区器是否有序是否默认是否易倾斜典型场景
HashPartitioner⚠️ 容易通用聚合
RangePartitioner✅ 较均匀排序、范围查询
CustomPartitioner自定义✅ 可控数据倾斜、业务规则

三、分区器生效条件(非常重要)

⚠️ 不是所有 RDD 都会使用 Partitioner

分区器只在以下情况生效:

  1. RDD 是 PairRDD(K, V)
  2. 使用 Shuffle 类算子
    • reduceByKey
    • groupByKey
    • join
    • cogroup
  3. 显式调用 partitionBy
rdd.reduceByKey(_ + _)        // 使用默认 HashPartitioner
rdd.partitionBy(new HashPartitioner(10))
    .reduceByKey(_ + _)      // 使用指定分区器

四、分区器与性能优化

✅ 减少 Shuffle

rdd.partitionBy(new HashPartitioner(10))
   .reduceByKey(_ + _)

👉 后续算子复用分区,避免重复 Shuffle

✅ 解决数据倾斜

  • 自定义 Partitioner
  • 将热点 key 单独拆分
  • 加盐(salting)+ 二次聚合

五、一句话总结

Spark 有三种分区器:HashPartitioner(默认)、RangePartitioner(有序)、CustomPartitioner(自定义)。
分区器决定 Shuffle 时 Key 的去向,直接影响数据分布、任务并行度和性能。

更多推荐