作家
登录

大数据系列之并行计算引擎Spark介绍

作者: 来源: 2017-04-24 14:44:07 阅读 我要评论

  • reduceByKey(func, [numTasks]) : 在一个(K,V)对的数据集上应用,返回一个(K,V)对的数据集,key雷同的值,都被应用指定的reduce函数聚合到一路。和groupbykey类似,义务的个数是可以经由过程第二个可选参数来设备的。
  • join(otherDataset, [numTasks]) :

在类型为(K,V)和(K,W)类型的数据集上调用,返回一个(K,(V,W))对,每个key中的所有元素都在一路的数据集

  • groupWith(otherDataset, [numTasks]) : 在类型为(K,V)和(K,W)类型的数据集上调用,返回一个数据集,构成元素为(K, Seq[V], Seq[W]) Tuples。这个操作在其它框架,称为CoGroup

计算两个石友间的合营石友数;

cartesian(otherDataset) : 笛卡尔积。但在数据集T和U上调用时,返回一个(T,U)对的数据集,所有元故旧互进行笛卡尔积。

  • flatMap(func) :

类似于map,然则每一个输入元素,会被映射为0到多个输出元素(是以,func函数的返回值是一个Seq,而不是单一元素)

Actions具体内容

  • reduce(func) : 经由过程函数func集合数据集中的所有元素。Func函数接收2个参数,返回一个值。这个函数必须是接洽关系性的,确保可以被精确的并发履行
  • collect() : 在Driver的法度榜样中,以数组的情势,返回数据集的所有元素。这平日会在应用filter或者其它操作后,返回一个足够小的数正人集再应用,直接将全部RDD集Collect返回,很可能会让Driver法度榜样OOM
  • count() : 返回数据集的元素个数
  • take(n) : 返回一个数组,由数据集的前n个元素构成。留意,这个操作今朝并非在多个节点上,并行履行,而是Driver法度榜样地点机械,单机计算所有的元素(Gateway的内存压力会增大年夜,须要谨慎应用)
  • first() : 返回数据集的第一个元素(类似于take(1))

saveAsTextFile(path) : 将数据集的元素,以textfile的情势,保存到本地文件体系,hdfs或者任何其它hadoop支撑的文件体系。Spark将会调用每个元素的toString办法,并将它转换为文件中的一行文本

  • saveAsSequenceFile(path) : 将数据集的元素,以sequencefile的格局,保存到指定的目次下,本地体系,hdfs或者任何其它hadoop支撑的文件体系。RDD的元素必须由key-value对构成,并都实现了Hadoop的Writable接口,或隐式可以转换为Writable(Spark包含了根本类型的转换,例如Int,Double,String等等)
  • foreach(func) : 在数据集的每一个元素上,运行函数func。这平日用于更新一?累加器变量,或者和外部存储体系做交互

算子分类

大年夜致可以分为三大年夜类算子:

  • Value数据类型的Transformation算子,这种变换并不触发提交功课,针对处理的数据项是Value型的数据。
  • Key-Value数据类型的Transfromation算子,这种变换并不触发提交功课,针对处理的数据项是Key-Value型的数据对。
  • Action算子,这类算子会触发SparkContext提交Job功课。

Spark在企业中的应用处景

  • 基于日记数据的快速萌芽体系营业;

  • Spark RDD cache/persist

Spark RDD cache

2.供给了多种缓存级别,以便于用户根据实际需求进行调剂

2.Spark>

3.cache应用

  • 之前用MapReduce实现过WordCount,如今我们用Scala实现下wordCount.是不是很简洁呢?!
  1. import org.apache.spark.{SparkConf, SparkContext} 
  2.  
  3. object SparkWordCount{ 
  4.  def main(args: Array[String]) { 
  5.  if (args.length == 0) { 
  6.  System.err.println("Usage: SparkWordCount <inputfile> <outputfile>"
  7.  System.exit(1) 
  8.  } 
  9.  
  10.  val conf = new SparkConf().setAppName("SparkWordCount"
  11.  val sc = new SparkContext(conf) 
  12.  
  13.  val file=sc.textFile("file:///hadoopLearning/spark-1.5.1-bin-hadoop2.4/README.md"
  14.  val counts=file.flatMap(line=>line.split(" ")) 
  15.  .map(word=>(word,1)) 
  16.  .reduceByKey(_+_) 

      推荐阅读

      数字化为丹东智慧城市腾飞安上羽翼

    第四,要让聪明城市的公平易近经由过程一张“市平易近卡”,融入社会生活并进行社会治理。市平易近卡将是成都聪明城市市平易近为行政办事体系供给同一的身份辨认模式。这张卡功能广泛,集传统的城市通卡、>>>详细阅读


    本文标题:大数据系列之并行计算引擎Spark介绍

    地址:http://www.17bianji.com/lsqh/34921.html

关键词: 探索发现

乐购科技部分新闻及文章转载自互联网,供读者交流和学习,若有涉及作者版权等问题请及时与我们联系,以便更正、删除或按规定办理。感谢所有提供资讯的网站,欢迎各类媒体与乐购科技进行文章共享合作。

网友点评
自媒体专栏

评论

热度

精彩导读
栏目ID=71的表不存在(操作类型=0)