spark shuffle解析
目录
shuflle
spark按照shuffle对stage进行了划分,有点类似hive,但是又有很大区别
先来看两个概念:
在划分stage时,最后一个stage称为finalStage,它本质上是一个ResultStage对象,前面的所有stage被称为ShuffleMapStage,每一个ShuffleMapStage的结束伴随着shuffle文件的写磁盘,而ResultStage意味着一个job的运行结束(即action算子)
发生了shuffle即代表数据要从一个task按照分区方式传输到另外task,这个过程就会有shuffle write 和shuffle reade
- shuffle write : 发生在shuf/fle之前,把需要shuffle的数据写到磁盘,保证了数据的可靠性
为什么在发送shuffle的时候,需要把数据保存到磁盘:- 避免占用内存太大出现oom(内存溢出)
- 保存到磁盘可以保证数据的安全性
- shuffle reade: 发生在shuffle之后,下游RDD需要读取上游RDD的数据
常见的shuffle算子
主动重分区
repartition :repartition只是coalesce接口中shuffle为true的简易实现,调用的还是coalesce算子
coalesce(true)
bykey自动重分区
reduceByKey
groupByKey
sortByKey
aggregateByKey
combineByKey
distinct() 是当key与value都一样的时候,会被当做重复的数据,需要对key操作
join重分区
join:两个RDD的分区数相同,在join的时候设置的分区数也相同,则在join阶段不会产生shuffle
具体是否发生shuffle还是看数据的交互,shuffle一定要发生宽依赖
- 比如我们要修改分区数,3个分区变成4个分区,肯定会shuffle,但是4个缩2个就不一定了
- 聚合操作,原先的key处于不同分区,聚合变成一个key,一个分区的数据会跑到不同分区里,shuffle
- 两个RDD的分区数相同,在join的时候设置的分区数也相同,则在join阶段不会产生shuffle
- 也就是说shuffle就是是有数据的分区会改变,分区的数据会打乱到多个分区里
- 广播变量的join,广播变量在所有节点的内存里,不涉及到分区,自然不会shuffle,但你要说我要在join时改变分区数,那也是会shuffle的,但是应该没有人会这么做
ShuffleManager
HashShuffleManager:
会产生大量的中间磁盘文件,进而由大量的磁盘IO操作影响了性能,在2.x版本被弃用了
未优化:
每一个 ShuffleMapTask 都会为每一个 ReducerTask 创建一个单独的文件,文件数就是map的数量*reduce的数量
shuffle read的拉取过程是一边拉取一边进行聚合的,将上一个stage的计算结果中的所有相同key,从各个节点上通过网络都拉取到自己所在的节点上,然后进行key的聚合或连接等操作
优化后
spark.shuffle.consolidateFiles设置为true
consolidate机制允许在一个 core 上连续执行的 ShuffleMapTasks 可以共用一个输出文件,这样就可以有效将多个task的磁盘文件进行一定程度上的合并,从而大幅度减少磁盘文件的数量,进而提升shuffle write的性能
每个Executor创建的磁盘文件的数量的计算公式为:CPU core的数量 * 下一个stage的task数量
SortShuffleManager
writer 有三种运行实现SortShuffleWriter,BypassMergeSortShuffleWriter,UnsafeShuffleWriter
reader 只有一种实现 BlockStoreShuffleReader
SortShuffleWriter
主要代码在这个类


主要过程:
- 数据会先写入一个内存数据结构中
根据不同的shuffle算子,可能选用不同的数据结构。如果是reduceByKey这种聚合类的shuffle算子,那么会选用Map数据结构,一边通过Map进行聚合,一边写入内存;如果是join这种普通的shuffle算子,那么会选用Array数据结构,直接写入内存 - 每写一条数据进入内存数据结构之后,就会判断一下,是否达到了某个临界阈值。如果达到临界阈值的话,那么就会尝试将内存数据结构中的数据溢写到磁盘,然后清空内存数据结构
- 在溢写到磁盘文件之前,会先根据key对内存数据结构中已有的数据进行排序。排序过后,会分批将数据写入磁盘文件。默认的batch数量是10000条
- 写入磁盘文件是通过Java的BufferedOutputStream实现的。BufferedOutputStream是Java的缓冲输出流,首先会将数据缓冲在内存中,当内存缓冲满溢之后再追加到该分区对应的文件中,这样可以减少磁盘IO次数,提升性能
- 最后会将之前所有的临时磁盘文件都进行合并,这就是merge过程,此时会将之前所有临时磁盘文件中的数据读取出来,然后依次写入最终的磁盘文件之中
- 由于一个task就只对应一个磁盘文件,也就意味着该task为下游stage的task准备的数据都在这一个文件中,因此还会单独写一份索引文件,其中标识了下游各个task的数据在文件中的start offset与end offset

BypassMergeSortShuffleWriter
- shuffle reduce task(即partition)数量小于
spark.shuffle.sort.bypassMergeThreshold参数的值,默认是200。 - 没有map side aggregations。也就是map端没有聚合操作,groupByKey 和combineByKey, 如果设定 mapSideCombine 为false,map端不会聚合
和SortShuffleWriter:
- 磁盘写机制不同
- 不会进行排序
- 也就是说,启用 BypassMerge 机制的最大好处在于,shuffle write 过程中,不需要进行数据的排序操作,也就节省掉了这部分的性能开销

Spark Shuffle 中数据结构
SortShuffleWriter 中使用 ExternalSorter 来对内存中的数据进行排序,ExternalSorter 中缓存记录数据的数据结构有两种,两者都是使用了 hash table 数据结构
PartitionedPairBuffer
mapSideCombine=false 时会使用该结构,也就是Map 阶段不进行 Combine 操作,是一个 Append-only Buffer,也就是仅支持向 Buffer 中追加数据键值对记录
PartitionedAppendOnlyMap
设置mapSideCombine=true时会使用该结构
更多推荐
所有评论(0)