直接看代码,照样有些晦涩难解,我们举个例子,一步一步解释锫:

按照膳绫擎的算法流程,大年夜致可以懂得:
- 抽样-->肯定界线(排序)
起首对spark有必定懂得的都应当知道,在spark中每个RDD可以懂得为一组分区,这些分区对应了内存块block,他们才是数据最终的载体。那么一个RDD由不合的分区构成,如许在处理一些map,filter等算子的时刻,就可以直接以分区为单位并行计算了。直到碰到shuffle的时刻才须要和其他的RDD合营。
在膳绫擎的图中,如不雅我们不特别设置的话,一个RDD由3个分区构成,那么在对它进行groupbykey的时刻,就会按照3进行分区。
然则如不雅是底层数据的问题,无论怎么竽暌古化,照样无法解决数据倾斜的。
- val sampleSize = math.min(20.0 * partitions, 1e6)
即采样数为60,每个分区取60个数。然则推敲到数据倾斜的情况,有的分区可能数据很多,是以在实际的采样时,会按照3倍大年夜小采样:
- val sampleSizePerPartition = math.ceil(3.0 * sampleSize / rdd.partitions.size).toInt
也就是说,最多会取60个样本数据。
然后就是遍历每个分区,取对应的样本数。
- val sketched = rdd.mapPartitionsWithIndex { (idx, iter) =>
- val seed = byteswap32(idx ^ (shift << 16))
- val (sample, n) = SamplingUtils.reservoirSampleAndCount(
- iter, sampleSizePerPartition, seed)
- //包装成三元组,(索引号,分区的内容个数,抽样的内容)
- Iterator((idx, n, sample))
- }.collect()
然后检查,是否有分区的样本数过多,如不雅多于平均值,则持续采样,这时直接用sample 就可以了
- sketched.foreach { case (idx, n, sample) =>
- if (fraction * n > sampleSizePerPartition) {
- imbalancedPartitions += idx
- }
推荐阅读
跟着数字化企业尽力寻求最佳安然解决筹划来保护其赓续扩大的收集,很多企颐魅正在寻求供给互操作性功能的下一代对象。软件定义收集(SDN)具有很多的优势,经由过程将多个设备的┞菲握平面整>>>详细阅读
本文标题:Spark源码分析之分区器的作用
地址:http://www.17bianji.com/lsqh/34919.html
1/2 1

网友点评
精彩导读
科技快报
品牌展示