Spark有几种分区器
·
在 Spark 中,**分区器(Partitioner)**决定了 RDD 中每条 Key 会被分配到哪一个分区,主要影响 groupByKey、reduceByKey、join 等 Shuffle 类算子 的行为。
Spark 内置的分区器主要有 三种,另外还支持 自定义分区器。
一、Spark 内置的三种分区器
1️⃣ HashPartitioner(默认)
类: org.apache.spark.HashPartitioner
核心原理
partitionId = key.hashCode() % numPartitions
特点
- 最常用、默认分区器
- 相同 key 一定进入同一个分区
- 实现简单、开销小
- 可能导致数据倾斜(hash 分布不均)
使用场景
reduceByKey、groupByKey、distinct- 不关心 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
核心原理
- 先对 key 进行采样
- 按 key 的大小排序
- 将 key 范围划分成
numPartitions个区间 - 每个区间一个分区
特点
- key 是有序的
- 尽量保证 每个分区数据量均匀
- 会产生 额外采样开销
- 适合排序、范围查询
使用场景
sortByKeyrepartitionAndSortWithinPartitions- 需要全局有序或范围查询
示例
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
分区器只在以下情况生效:
- RDD 是 PairRDD(K, V)
- 使用 Shuffle 类算子
reduceByKeygroupByKeyjoincogroup
- 显式调用
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 的去向,直接影响数据分布、任务并行度和性能。
更多推荐
所有评论(0)