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

spark core源码分析16 Shuffle详解-读流程,sparkshuffle

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

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

spark core源码分析16 Shuffle详解-读流程,sparkshuffle


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


shuffle的读流程也是从compute方法开始的

override def compute(split: Partition, context: TaskContext): Iterator[(K, C)] = {
    val dep = dependencies.head.asInstanceOf[ShuffleDependency[K, V, C]]
    SparkEnv.get.shuffleManager.getReader(dep.shuffleHandle, split.index, split.index + 1, context)
      .read()
      .asInstanceOf[Iterator[(K, C)]]
  }

目前来说,不管是sortShuffleManager还是hashShuffleManager,getReader方法返回的都是HashShuffleReader。

接着调用read方法,如下:

/** Read the combined key-values for this reduce task */
  override def read(): Iterator[Product2[K, C]] = {
    val ser = Serializer.getSerializer(dep.serializer)
    val iter = BlockStoreShuffleFetcher.fetch(handle.shuffleId, startPartition, context, ser)

    val aggregatedIter: Iterator[Product2[K, C]] = if (dep.aggregator.isDefined) {
      if (dep.mapSideCombine) {
        new InterruptibleIterator(context, dep.aggregator.get.combineCombinersByKey(iter, context))
      } else {
        new InterruptibleIterator(context, dep.aggregator.get.combineValuesByKey(iter, context))
      }
    } else {
      require(!dep.mapSideCombine, "Map-side combine without Aggregator specified!")

      // Convert the Product2s to pairs since this is what downstream RDDs currently expect
      iter.asInstanceOf[Iterator[Product2[K, C]]].map(pair => (pair._1, pair._2))
    }

    // Sort the output if there is a sort ordering defined.
    dep.keyOrdering match {
      case Some(keyOrd: Ordering[K]) =>
        // Create an ExternalSorter to sort the data. Note that if spark.shuffle.spill is disabled,
        // the ExternalSorter won't spill to disk.
        val sorter = new ExternalSorter[K, C, C](ordering = Some(keyOrd), serializer = Some(ser))
        sorter.insertAll(aggregatedIter)
        context.taskMetrics.incMemoryBytesSpilled(sorter.memoryBytesSpilled)
        context.taskMetrics.incDiskBytesSpilled(sorter.diskBytesSpilled)
        sorter.iterator
      case None =>
        aggregatedIter
    }
  }

该方法首先调用了fetch方法,介绍一下

1、在task运行那节介绍过,shuffleMapTask运行完成后,会将shuffleId及mapstatus的映射注册到mapOutputTracker中

2、fetch方法首先尝试在本地mapstatuses中查找是否有该shuffleId的信息,有则本地取;否则想master的mapOutputTracker请求并读取,返回块管理器的地址和对应partition的文件长度

3、然后根据我们得到的shuffleId等信息去remote或者local通过netty/nio读取,返回一个迭代器

4、返回的迭代器中的数据并不是全部在内存中的,读取时会根据配置的内存最大值来读取。内存不够的话,下一个待读取


fetch方法返回一个迭代器后,根据是否mapSideCombine来区分时候需要将读取到的数据进行合并操作。合并过程与写流程类似,内存放不下就写入本地磁盘。

如果还需要keyOrdering的,new一个ExternalSorter进行外部排序。之后也是同shuffle写流程的insertAll。


版权声明:本文为博主原创文章,未经博主允许不得转载。

分享到: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

相关文章

  • RPC原理,rpc
  • Elasticsearch 之 Facet,elasticsearchfacet
  • Mellanox网卡,使用netperf进行性能测试,mellanoxnetperf
  • Spark中广播变量知识点
  • 运行spark-shell时遇到的主机地址的错误,spark-shell主机
  • 我与英语技术书籍,英语技术书籍
  • 别再用高考来绑架孩子们心中的“公平”,孩子们心中
  • spark core源码分析8 从简单例子看transformation,transformation
  • hive:Access denied for user 'root'@'%',hivedenied
  • Hbase Client Test Case,hbasecase

文章分类

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

最近更新的内容

    • 常见类型网站的SEO应该怎么操作?类型不同,操作自然不同,因为分门别类嘛,seo分门别类
    • hadoop-2.6.0安装(new),hadoop-2.6.0new
    • 个性化推荐的十大挑战,个性化十大挑战
    • hadoop-2.6.0 Unhealthy Nodes 问题,hadoop2.6.0安装
    • Hadoop源码分析----RPC反射机制,hadoop----rpc
    • MapReduce中Shuffle过程整理,mapreduceshuffle
    • 每日定时导入hive数据仓库的自动化脚本,hive数据仓库脚本
    • 创建GZIP压缩格式的HIVE表,gzip压缩格式hive
    • 机器学习-Python中训练模型的保存和再使用,-python模型
    • zookeeper的安装,搭建,zookeeper安装搭建

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

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