• linkedu视频
  • 平面设计
  • 电脑入门
  • 操作系统
  • 办公应用
  • 电脑硬件
  • 动画设计
  • 3D设计
  • 网页设计
  • CAD设计
  • 影音处理
  • 数据库
  • 程序设计
  • 认证考试
  • 信息管理
  • 信息安全
菜单
linkedu.com
导航菜单
  • 网页制作
  • 数据库
  • 程序设计
  • 操作系统
  • CMS教程
  • 游戏攻略
  • 脚本语言
  • 平面设计
  • 软件教程
  • 网络安全
  • 电脑知识
  • 服务器
  • 视频教程
  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx
您的位置:首页 > 服务器 >云计算 > spark core源码分析15 Shuffle详解-写流程,sparkshuffle

spark core源码分析15 Shuffle详解-写流程,sparkshuffle

作者:网友 字体:[增加 减小] 来源:互联网

本文主要包含spark core 2.10,spark shuffle,spark中的shuffle,spark的shuffle过程,spark shuffle算子等服务器相关知识,网友希望可以进行参考

spark core源码分析15 Shuffle详解-写流程,sparkshuffle



博客地址: http://blog.csdn.net/yueqian_zhu/


Shuffle是一个比较复杂的过程,有必要详细剖析一下内部写的逻辑

ShuffleManager分为SortShuffleManager和HashShuffleManager

一、SortShuffleManager

每个ShuffleMapTask不会为每个Reducer生成一个单独的文件;相反,它会将所有的结果写到一个本地文件里,同时会生成一个index文件,Reducer可以通过这个index文件取得它需要处理的数据。避免产生大量的文件的直接收益就是节省了内存的使用和顺序Disk IO带来的低延时。

它在写入分区数据的时候,首先会根据实际情况对数据采用不同的方式进行排序操作,底线是至少按照Reduce分区Partition进行排序,这样同一个Map任务Shuffle到不同的Reduce分区中去的所有数据都可以写入到同一个外部磁盘文件中去,用简单的Offset标志不同Reduce分区的数据在这个文件中的偏移量。这样一个Map任务就只需要生成一个shuffle文件,从而避免了上述HashShuffleManager可能遇到的文件数量巨大的问题


/** Get a writer for a given partition. Called on executors by map tasks. */
  override def getWriter[K, V](handle: ShuffleHandle, mapId: Int, context: TaskContext)
      : ShuffleWriter[K, V] = {
    val baseShuffleHandle = handle.asInstanceOf[BaseShuffleHandle[K, V, _]]
    shuffleMapNumber.putIfAbsent(baseShuffleHandle.shuffleId, baseShuffleHandle.numMaps)
    new SortShuffleWriter(
      shuffleBlockResolver, baseShuffleHandle, mapId, context)
  }

shuffleMapNumber是一个HashMap<shuffleId,numMaps>

SortShuffleWriter提供write接口用于真实数据的写磁盘,而在write接口中会使用shuffleBlockResolver与底层文件打交道


下面看获得SortShuffleWriter之后,调用write进行写

writer.write(rdd.iterator(partition, context).asInstanceOf[Iterator[_ <: Product2[Any, Any]]])
write的参数其实就是调用了rdd的compute方法进行计算,返回的这个partition的迭代器
/** Write a bunch of records to this task's output */
  override def write(records: Iterator[Product2[K, V]]): Unit = {
    if (dep.mapSideCombine) {
      require(dep.aggregator.isDefined, "Map-side combine without Aggregator specified!")
      sorter = new ExternalSorter[K, V, C](
        dep.aggregator, Some(dep.partitioner), dep.keyOrdering, dep.serializer)
      sorter.insertAll(records)
    } else {
      // In this case we pass neither an aggregator nor an ordering to the sorter, because we don't
      // care whether the keys get sorted in each partition; that will be done on the reduce side
      // if the operation being run is sortByKey.
      sorter = new ExternalSorter[K, V, V](None, Some(dep.partitioner), None, dep.serializer)
      sorter.insertAll(records)
    }

    // Don't bother including the time to open the merged output file in the shuffle write time,
    // because it just opens a single file, so is typically too fast to measure accurately
    // (see SPARK-3570).
    val outputFile = shuffleBlockResolver.getDataFile(dep.shuffleId, mapId)
    val blockId = ShuffleBlockId(dep.shuffleId, mapId, IndexShuffleBlockResolver.NOOP_REDUCE_ID)
    val partitionLengths = sorter.writePartitionedFile(blockId, context, outputFile)
    shuffleBlockResolver.writeIndexFile(dep.shuffleId, mapId, partitionLengths)

    mapStatus = MapStatus(blockManager.shuffleServerId, partitionLengths)
  }

可以看到,设置了mapSideCombine的需要将aggregator和keyOrdering传入到ExternalSorter中,否则将上面两项参数设为None。接着调用insertAll方法

def insertAll(records: Iterator[_ <: Product2[K, V]]): Unit = {
    // TODO: stop combining if we find that the reduction factor isn't high
    val shouldCombine = aggregator.isDefined

    if (shouldCombine) {
      // Combine values in-memory first using our AppendOnlyMap
      val mergeValue = aggregator.get.mergeValue
      val createCombiner = aggregator.get.createCombiner
      var kv: Product2[K, V] = null
      val update = (hadValue: Boolean, oldValue: C) => {
        if (hadValue) mergeValue(oldValue, kv._2) else createCombiner(kv._2)
      }
      while (records.hasNext) {
        addElementsRead()
        kv = records.next()
        map.changeValue((getPartition(kv._1), kv._1), update)
        maybeSpillCollection(usingMap = true)
      }
    } else if (bypassMergeSort) {
      // SPARK-4479: Also bypass buffering if merge sort is bypassed to avoid defensive copies
      if (records.hasNext) {
        spillToPartitionFiles(
          WritablePartitionedIterator.fromIterator(records.map { kv =>
            ((getPartition(kv._1), kv._1), kv._2.asInstanceOf[C])
          })
        )
      }
    } else {
      // Stick values into our buffer
      while (records.hasNext) {
        addElementsRead()
        val kv = records.next()
        buffer.insert(getPartition(kv._1), kv._1, kv._2.asInstanceOf[C])
        maybeSpillCollection(usingMap = false)
      }
    }
  }
解释一下内部逻辑:

(1)  如果是shouldCombine,将k-v信息记录到一个Array中,默认是大小是64*2,存储格式为key0,value0,key1,value1,key2,value2...。map.changeValue方法就是将value的值不断的调用mergeValue方法,去更新array中指定位置的value值。如果k-v的量达到array size的0.7时,会自动扩容。

之后调用maybeSpillCollection,首先判断是否需要spill,依据是开启spillingEnabled标志(不开启有OOM风险,其实上面rehash扩容的时候就应该是有OOM风险了),且读取的元素是32的整数倍,且目前占用的内存大于设置的阀值(5M),就去向shuffleMemoryManager申请内存(shuffleMemoryManager中有一个阀值,每个shuffle task向他申请时都会记录一下),申请的容量是当前使用容量的2倍减去阀值(5M),如果申请成功就增加阀值。如果目前内存占用量还是大于新的阀值,则必须要进行spill了,否则认为内存还够用。真正spill操作之后,释放刚才从shuffleMemoryManager中申请的内存以及还原阀值到初始值(5M)。

spill方法:如果partition数量<=200,且没有设置map端的combine,就调用spillToPartitionFiles方法,否则调用spillToMergeableFile方法,之后会讲到。

所以在这个分支而言,我们是shouldCombine的,所以调用的是spillToMergeFile方法。

需要注意的是,在spill之前,我们是有一个数据结构来保存数据的,有map和buffer可选择。由于shouldCombine是有可能去更新数据的,即调用我们的mergeValue方法之类的,所以我们用map。

(2)  如果是bypassMergeSort(partition数量<=200,且没有设置map端的combine),调用的是spillToPartitionFiles方法。这种模式直接写partition file,就没有缓存这一说了。

(3)  如果是 非shouldCombine,非bypassMergeSort,这里因为我们不需要merge操作,直接使用buffer作为spill前的缓存结构。之后调用maybeSpillCollection方法。


看一下spillToMergeableFile方法:

(1) 在localDirs下面的子目录下创建一个写shuffle的文件

(2) 对缓存中的数据进行排序,原则是按partitionID和partition内的key排序,得到的数据格式为((partitionId_0,key_0),value_0),((partitionId_0,key_1),value_1)......((partition

分享到:QQ空间新浪微博腾讯微博微信百度贴吧QQ好友复制网址打印

您可能想查找下面的文章:

  • spark core源码分析16 Shuffle详解-读流程,sparkshuffle
  • spark core源码分析15 Shuffle详解-写流程,sparkshuffle
  • spark core源码分析14 参数配置,sparkcore
  • spark core源码分析10 Task的运行,sparkcore
  • spark core源码分析9 从简单例子看action操作,sparkcore
  • spark core源码分析8 从简单例子看transformation,transformation

相关文章

  • HBase数据导入的几种操作,hbase数据导入几种
  • 在Ubuntu 14.04安装和使用Docker,14.04docker
  • 孙其功陪你学之——Spark 正则化和SparkSQL,孙其功sparksql
  • MapReduce的两表join操作优化,mapreduce表join
  • A Note on Distributed Computing,anoteondialectic
  • CSR1000V在XenServer的安装和简单使用,csr1000vxenserver
  • Mahout-HashMap的进化版FastByIdMap,mahout
  • Spark下的PageRank实现,SparkPageRank实现
  • spark资料下载,spark下载
  • 推荐引擎mahout安装与配置,引擎mahout配置

文章分类

  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx

最近更新的内容

    • ZooKeeper+Hadoop2.6.0的ResourceManager HA搭建,hadoop2.6集群搭建
    • Linux配置时间服务器,linux配置服务器
    • 可穿戴KEY带来的身份认证的革命,key身份认证
    • mongoVUE的增删改查操作使用说明;一、查询;1、精确查询;1)右键点击集合名,再左键点击Find;或者直接点击工具栏上的Find;2)查询界面,包括四个区域;{Find}区,查询条件格式{&quot;se,m
    • HBase shell 启动出错 org.apache.zookeeper.KeeperException$ConnectionLossException: KeeperErrorCode = Con,
    • 修改 openstack 中 nova boot 创建实例只能在10个以内的限制,openstacknova
    • HDP出现Could not create the Java Virtual Machine解决方法,hdpvirtual
    • Hive之简单查询不启用MapReduce,hive启用mapreduce
    • 大数据处理算法一:Bitmap算法,数据处理bitmap算法
    • Namespace:Openstack的网络实现,namespaceopenstack

关于我们 - 联系我们 - 免责声明 - 网站地图

©2020-2025 All Rights Reserved. linkedu.com 版权所有