Spark核心源码分析

01 Spark集群启动原理

Spark集群启动时,会在当前节点(脚本执行节点)上启动Master,在配置文件conf/slave中指定的每个节点上启动一个Worker。而Spark集群是通过脚本启动的。
查看启动脚本的内容:

$ cat sbin/start-all.sh

内容:

#!/usr/bin/env bash
#启动所有spark守护进程
#在此节点上启动Master
#在conf/slave中指定的每个节点上启动一个Worker
 #如果SPARK_HOME环境变量为空,则使用export设置
if [ -z "${SPARK_HOME}" ]; then
  export SPARK_HOME="$(cd "<code>dirname "$0""/..; pwd)"
fi
#加载Spark配置
. "${SPARK_HOME}/sbin/spark-config.sh"
#启动Master
"${SPARK_HOME}/sbin"/start-master.sh
#启动Workers
"${SPARK_HOME}/sbin"/start-slaves.sh

从启动脚本的内容和注释可以看出,start-all.sh脚本主要做了四件事:

(1)检查并设置环境变量
(2)加载Spark配置
Spark配置的加载,使用了脚本文件sbin/spark-config.sh。
(3)启动Master进程
Master进程的启动,使用了脚本文件sbin/start-master.sh。
(4)启动Worker进程
Worker进程的启动,使用了脚本文件sbin/start-slaves.sh,该文件负责批量启动集群中各个节点上的Worker进程。

02 Spark应用程序提交原理

Spark应用程序的提交入口是spark-submit脚本。除此之外,也可以使用spark-shell脚本和spark-sql脚本进入交互式命令行执行Spark程序。
查看spark-shell脚本的内容:

$ cat bin/spark-shell

内容关键代码:

#!/usr/bin/env bash
......
function main() {
  if $cygwin; then
    stty -icanon min 1 -echo > /dev/null 2>&1
    export SPARK_SUBMIT_OPTS="$SPARK_SUBMIT_OPTS -Djline.terminal=unix"
"${SPARK_HOME}"/bin/spark-submit --class org.apache.spark.repl.Main --name
 "Spark shell" "$@"
    stty icanon echo > /dev/null 2>&1
  else
    export SPARK_SUBMIT_OPTS
"${SPARK_HOME}"/bin/spark-submit --class org.apache.spark.repl.Main --name
 "Spark shell" "$@"
  fi
}
......

从关键代码可以看出,spark-shell脚本调用了应用程序提交脚本spark-submit,并传入了org.apache.spark.repl.Main类和使用spark-shell命令时的所有参数(例如--master)。通过org.apache.spark.repl.Main类启动的spark-shell将进入REPL模式(交互式编程环境),一直等待用户的输入,除非手动退出。
查看spark-sql脚本的内容:

$ cat bin/spark-sql

代码如下:

#!/usr/bin/env bash
if [ -z "${SPARK_HOME}" ]; then
  source "$(dirname "$0")"/find-spark-home
fi

export _SPARK_CMD_USAGE="Usage: ./bin/spark-sql [options] [cli option]"
exec "${SPARK_HOME}"/bin/spark-submit --class 
org.apache.spark.sql.hive.thriftserver.SparkSQLCLIDriver "$@"

从上述代码可以看出,spark-sql脚本同样调用了应用程序提交脚本spark-submit。这说明,无论采用哪种方式提交应用程序,都间接的调用了spark-submit脚本。查看spark-submit脚本的内容,发现最后一行有如下代码:

exec "${SPARK_HOME}"/bin/spark-class org.apache.spark.deploy.SparkSubmit "$@“

可以看出,最后使用exec命令执行了脚本bin/spark-class,并传入了org.apache.spark.deploy.SparkSubmit类和使用spark-submit命令时的所有参数。脚本bin/spark-class用于执行Spark应用程序并启动JVM进程,而这里执行的正是SparkSubmit类。

从上面的分析可以总结出,若通过spark-shell进行启动,将默认传入以下参数:

org.apache.spark.deploy.SparkSubmit
--class org.apache.spark.repl.Main
--name "Spark shell"

若通过spark-submit提交应用程序,将默认传入以下参数:

org.apache.spark.deploy.SparkSubmit

03 Spark作业工作原理

在学习Spark作业的工作原理时,通常会引入Hadoop MapReduce的工作原理作为入门比较,因为MapReduce与Spark的工作原理有很多相似之处。

1、MapReduce工作原理

MapReduce计算模型主要由三个阶段组成:Map阶段、Shuffle阶段、Reduce阶段。

file

(1)Map阶段

将输入的多个分片(Split)由Map任务以完全并行的方式处理。每个分片由一个Map任务来处理。默认情况下,输入分片的大小与HDFS中数据块(Block)的大小是相同的,即文件有多少个数据块就有多少个输入分片,也就会有多少个Map任务。从而可以调整HDFS数据块的大小来间接改变Map任务的数量。

每个Map任务对输入分片中的记录按照一定的规则解析成多个<key,value>对。默认将文件中的每一行文本内容解析成一个<key,value>对,key为每一行的起始位置,value为本行的文本内容。然后将解析出的所有<key,value>对分别输入到map()方法中进行处理(map()方法一次只处理一个<key,value>对)。map()方法将处理结果仍然是以<key,value>对的形式进行输出。

在数据溢写到磁盘之前,会对数据进行分区(Partition)。分区的数量与设置的Reduce任务的数量相同(默认Reduce任务的数量为1,可以在编写MapReduce程序时对其修改)。这样每个Reduce任务会处理一个分区的数据,可以防止有的Reduce任务分配的数据量太大,而有的Reduce任务分配的数据量太小,从而可以负载均衡,避免数据倾斜。数据分区的划分规则为:取<key,value>对中key的hashCode值,然后除以Reduce任务数量后取余数,余数则是分区编号,分区编号一致的<key,value>对则属于同一个分区。因此,key值相同的<key,value>对一定属于同一个分区,但是同一个分区中可能有多个key值不同的<key,value>对。由于默认Reduce任务的数量为1,而任何数字除以1的余数总是0,因此分区编号从0开始。

(2)Reduce阶段

Reduce阶段首先会对Map阶段的输出结果按照分区进行再一次合并,将同一分区的<key,value>对合并到一起,然后按照key对分区中的<key,value>对进行排序。
每个分区会将排序后的<key,value>对按照key进行分组,key相同的<key,value>对将合并为<key,value-list>对,最终每个分区形成多个<key,value-list>对。例如,key中存储的是用户ID,则同一个用户的<key,value>对会合并到一起。
排序并分组后的分区数据会输入到reduce()方法中进行处理,reduce()方法一次只能处理一个<key,value-list>对。
最后,reduce()方法将处理结果仍然以<key,value>对的形式通过context.write(key,value)进行输出。

(3)Shuffle阶段

Shuffle阶段所处的位置是Map任务输出后,Reduce任务接收前。主要是将Map任务的无规则输出形成一定的有规则数据,以便Reduce任务进行处理。
总结来说,MapReduce的工作原理主要是:通过Map任务读取HDFS中的数据块,这些数据块由Map任务以完全并行的方式处理;然后将Map任务的输出进行排序后输入到Reduce任务中;最后Reduce任务将计算的结果输出到HDFS文件系统中。

2、Spark工作原理

典型的Spark作业(Job)的工作流程如图:

file

(1)从数据源(本地文件、HDFS、HBase等)读取数据并创建RDD。
(2)对RDD进行一系列的转化操作。
(3)对最终RDD执行行动操作,开始一系列的计算,产生计算结果。
(4)将计算结果发送到Driver端,进行查看和输出。

04 Spark检查点原理

通过调用SparkContext的setCheckpointDir()方法指定检查点数据的存储路径,即把RDD数据放到什么位置,通常是放到HDFS中(如果在集群中运行,则必须是HDFS目录)。

总结来说,RDD检查点的运行流程如下:
(1)通过SparkContext的setCheckpointDir()方法设置RDD检查点数据的存储路径。
(2)调用RDD的checkpoint()方法将RDD标记为检查点。
(3)当在RDD上运行一个作业后,会立即触发RDDCheckpointData中的checkpoint()方法。
(4)在checkpoint()方法中又执行了doCheckpoint()方法。
(5)doCheckpoint()方法中执行了writeRDDToCheckpointDirectory()方法。
(6)writeRDDToCheckpointDirectory()方法内部通过调用runJob()方法运行一个作业,真正将RDD数据写入到检查点目录中,写入完成后返回一个ReliableCheckpointRDD实例。

Views: 105

Spark RDD弹性分布式数据集

01 什么是RDD

Spark提供了一种对数据的核心抽象,称为弹性分布式数据集(Resilient Distributed Dataset,简称RDD)。这个数据集的全部或部分可以缓存在内存中,并且可以在多次计算时重用。RDD其实就是一个分布在多个节点上的数据集合。

RDD的弹性主要是指:当内存不够时,数据可以持久化到磁盘,并且RDD具有高效的容错能力。分布式数据集是指:一个数据集存储在不同的节点上,每个节点存储数据集的一部分。

例如,将数据集(hello,world,scala,spark,love,spark,happy)存储在三个节点上,节点一存储(hello,world),节点二存储(scala,spark,love),节点三存储(spark,happy),这样对三个节点的数据可以并行计算,并且三个节点的数据共同组成了一个RDD。

file

分布式数据集类似于HDFS中的文件分块,不同的块存储在不同的节点上;而并行计算类似使用MapReduce读取HDFS中的数据并进行Map和Reduce操作。Spark则包含了这两种功能,并且计算更加灵活。

RDD的主要特征如下:

  • RDD是不可变的,但是可以将RDD转换成新的RDD进行操作
  • RDD是可分区的,一个RDD可以由很多分区组成(数据存在多个节点上),每个分区对应一个Task来执行
  • 对RDD的操作是作用在RDD的每个分区上
  • RDD拥有一系列对分区进行计算的函数,成为算子
  • RDD之间存在依赖关系,可以实现管道化,避免中间数据的存储

在编程时,可以把RDD看作是一个数据操作的基本单位,而不必关心数据的分布式特性,Spark会自动将RDD的数据分发到集群的各个节点。Spark中对数据的操作主要是对RDD的操作(创建、转化、求值)。

1、从对象集合创建RDD

Spark可以通过parallelize()或makeRDD()方法将一个对象集合转化为RDD。
例如,将一个List集合转化为RDD:

scala> val rdd=sc.parallelize(List(1,2,3,4,5,6))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0]

scala> val rdd=sc.makeRDD(List(1,2,3,4,5,6))
rdd: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[1]

从返回信息可以看出,上述创建的RDD中存储的是Int类型数据。实际上RDD也是一个集合(并行集合),与常用的List集合不同的是,RDD集合的数据分布于多台机器上。

2、从外部存储创建RDD

Spark的textFile()方法可以读取本地文件系统或外部其它系统中的数据,并创建RDD。所不同的是,数据的来源路径不同。

首先在本地准备文件"/tmp/words.txt,内容如下:

hello hadoop
hello java
scala

进入Spark Shell读取本地数据(加载本地文件,必须以file:///开头):

//读取本地数据
scala> val rdd=sc.textFile("file:///tmp/words.txt")
rdd: org.apache.spark.rdd.RDD[String] = /tmp/words.txt MapPartitionsRDD[1]

scala> rdd.collect
res1: Array[String] = Array("hello hadoop ", "hello java ", "scala ")

异常问题

Spark Shell 调用 rdd.collect报错

spark Compression codec com.hadoop.compression.lzo.LzoCodec not found.

这是因为在hadoop中配置了编解码器lzo,所以当使用yarn模式时,spark自身没有lzo的jar包所以报错无法找到!
解决办法:
配置spark-default.conf文件!
修改spark-defaults.sh, 添加:

spark.jars /opt/pkg/hadoop/share/hadoop/common/hadoop-lzo-0.4.20.jar
读取本地数据时报错路径不存在

相似问题

这个问题发生在集群方式提交本地文件的时候。

这是因为文件的读取是发生在Executor所在的节点(即Worker节点),只有将文件tmp/words.txt同步到所有Worker节点才能避免路径不存在的问题发生。

有意思的地方是,即使不同步到所有Worker节点,先调用rdd.first拉取第一条数据之后再调用rdd.collect就可以获取到数据了(中间也会出现文件找不到问题,但是仍可以获取到结果)。原因就是RDD的collect方法属于转化算子(算子的概念后面会讲)是有惰性的(不是一开始马上收集,需要等待动作算子(如first或者count)操作数据时才会收集)才会主动收集数据

解决办法就是:

  1. 改为Local模式运行 $ spark-shell --master local
  2. 把文件file:///tmp/words.txt同步到所有salve节点。
  3. 把文件先上传到HDFS,然后再读取HDFS上的文件

读取HDFS数据

先将数据上传到HDFS

$ hadoop fs -put /tmp/words.txt /tmp/

再读取HDFS数据

//读取HDFS数据
scala> val rdd=sc.textFile("hdfs://centos01:8082/tmp/words.txt")
rdd: org.apache.spark.rdd.RDD[String] = hdfs://centos01:8082/tmp/words.txt MapPartitionsRDD[2]

scala> rdd.collect
res2: Array[String] = Array("hello hadoop ", "hello java ", "scala ")

02 RDD常用算子

RDD被创建后是只读的,不允许修改。Spark提供了丰富的用于操作RDD的方法,这些方法被称为算子。一个创建完成的RDD只支持两种算子:转化(Transformation)算子和行动(Action)算子。

1、转化算子

Spark中的转化算子主要时对RDD进行操作并返回新的RDD。

(1)map(func)

对rdd1应用map()算子,将rdd1中的每个元素加1并返回一个名为rdd2的新RDD:

scala> val rdd1=sc.parallelize(List(1,2,3,4,5,6))
scala> val rdd2=rdd1.map(x => x+1)

(2)filter(func)

对源RDD的每个元素应用函数func进行过滤,并返回一个新的RDD。
例如,过滤出rdd1中大于3的所有元素,并输出结果:

scala> val rdd1=sc.parallelize(List(1,2,3,4,5,6))
scala> val rdd2=rdd1.filter(_>3)
scala> rdd2.collect
res1: Array[Int] = Array(4, 5, 6)

上述代码中的下划线“_”代表rdd1中的每个元素。

(3)flatMap(func)

与map()算子类似,但是每个传入给函数func的RDD元素会返回0到多个元素,最终会将返回的所有元素合并到一个RDD。
例如,将集合List转为rdd1,然后调用rdd1的flatMap()算子将rdd1的每个元素按照空格进行分割成多个元素,最终合并所有元素到一个新的RDD。

scala> val rdd1=sc.parallelize(List("hadoop hello scala","spark hello"))
scala> val rdd2=rdd1.flatMap(_.split(" "))
scala> rdd2.collect
res3: Array[String] = Array(hadoop, hello, scala, spark, hello)

file

(4)reduceByKey()

reduceByKey()算子的作用对像是元素为(key,value)形式(Scala元组)的RDD,使用该算子可以将相同key的元素聚集到一起,最终把所有相同key的元素合并成为一个元素。该元素的key不变,value可以聚合成一个列表或者进行求和等操作。最终返回的RDD的元素类型和原有类型保持一致。
例如,有两个同学zhangsan和lisi,zhangsan的语文和数学成绩分别为98、78,lisi的语文和数学成绩分别为88、79,现需要分别求zhangsan和lisi的总成绩,代码如下:

scala> val list=List(("zhangsan",98),("zhangsan",78),("lisi",88),("lisi",79))
scala> val rdd1=sc.parallelize(list)
scala> val rdd2=rdd1.reduceByKey((x,y)=>x+y) //可以简化为rdd1.reduceByKey(_+_)
scala> rdd2.collect
res5: Array[(String, Int)] = Array((zhangsan,176), (lisi,167))

file

(5)groupByKey()

groupByKey()算子的作用对像是元素为(key,value)形式(Scala元组)的RDD,使用该算子可以将相同key的元素聚集到一起,最终把所有相同key的元素合并成为一个元素。该元素的key不变,value则聚集到一个集合中。
仍然以上述求学生zhangsan和lisi的总成绩为例,使用groupByKey()算子的代码如下:

scala> val list=List(("zhangsan",98),("zhangsan",78),("lisi",88),("lisi",79))
scala> val rdd1=sc.parallelize(list)
scala> val rdd2=rdd1.groupByKey()
scala> rdd2.map(x => (x._1,x._2.sum)).collect
res0: Array[(String, Int)] = Array((zhangsan,176), (lisi,167))

file

(6)union()

union()算子将两个RDD合并为一个新的RDD,主要用于对不同的数据来源进行合并,两个RDD中的数据类型要保持一致。
例如以下代码,通过集合创建了两个RDD,然后将两个RDD合并成了一个RDD:

scala> val rdd1=sc.parallelize(Array(1,2,3))
scala> val rdd2=sc.parallelize(Array(4,5,6))
scala> val rdd3=rdd1.union(rdd2)
scala> rdd3.collect
res8: Array[Int] = Array(1, 2, 3, 4, 5, 6) 

(7)sortBy()

sortBy()算子将RDD中的元素按照某个规则进行排序。该算子的第一个参数为排序函数,第二个参数是一个布尔值,指定升序(默认)或降序。若需要降序排列,需将第二个参数置为false。
例如,一个数组中存放了三个元组,将该数组转为RDD集合,然后对该RDD按照每个元素中的第二个值进行降序排列,代码如下:

scala> val rdd1=sc.parallelize(Array(("hadoop",12),("java",32),("spark",22)))
scala> val rdd2=rdd1.sortBy(x=>x._2,false)
scala> rdd2.collect
res2: Array[(String, Int)] = Array((java,32),(spark,22),(hadoop,12))

(8)sortByKey()

sortByKey()算子将(key,value)形式的RDD按照key进行排序。默认升序,若需降序排列,可以传入参数false,代码如下:
rdd.sortByKey(false)

(9)join ()

join()算子将两个(key,value)形式的RDD根据key进行连接操作,相当于数据库的内连接(inner join),只返回两个RDD都匹配的内容。例如,将rdd1和rdd2进行内连接,代码如下:

scala> val arr1= Array(("A","a1"),("B","b1"),("C","c1"),("D","d1"),("E","e1")
scala> val rdd1 = sc.parallelize(arr1))
scala> val arr2= Array(("A","A1"),("B","B1"),("C","C1"),("C","C2"),("C","C3"),("E","E1"))
scala> val rdd2 = sc.parallelize(arr2)
scala> rdd1.join(rdd2).collect
res0: Array[(String, (String, String))] = Array((B,(b1,B1)), (A,(a1,A1)), (C,(c1,C1)), (C,(c1,C2)), (C,(c1,C3)), (E,(e1,E1)))
scala> rdd2.join(rdd1).collect
res1: Array[(String, (String, String))] = Array((B,(B1,b1)), (A,(A1,a1)), (C,(C1,c1)), (C,(C2,c1)), (C,(C3,c1)), (E,(E1,e1)))

file

2、行动算子

Spark中的转化算子并不会马上进行运算,而是在遇到行动算子时才会执行相应的语句,触发Spark的任务调度。

file

2、行动算子

(1)reduce()

将数字1到100所组成的集合转为RDD,然后对该RDD进行reduce()算子计算,统计RDD中所有元素值的总和,代码如下:

scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.reduce(_+_)
res2: Int = 5050
scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.count
res3: Long = 100

(2)count()

统计RDD集合中元素的数量:

scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.reduce(_+_)
res2: Int = 5050
scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.count
res3: Long = 100

(3)countByKey()

例如,List集合中存储的是键值对形式的元组,使用该List集合创建一个RDD,然后对其进行countByKey()的计算:

scala> val rdd1 = sc.parallelize(List(("zhang",87),("zhang",79),("li",90)))
scala> rdd1.countByKey
res1: scala.collection.Map[String,Long] = Map(zhang -> 2, li -> 1) 

scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.take(5)
res4: Array[Int] = Array(1, 2, 3, 4, 5)

(4)take(n)

返回集合中前5个元素组成的数组:

scala> val rdd1 = sc.parallelize(List(("zhang",87),("zhang",79),("li",90)))
scala> rdd1.countByKey
res1: scala.collection.Map[String,Long] = Map(zhang -> 2, li -> 1) 

scala> val rdd1 = sc.parallelize(1 to 100)
scala> rdd1.take(5)
res4: Array[Int] = Array(1, 2, 3, 4, 5)

03 RDD的分区

RDD是一个大的数据集合,该集合被划分成多个子集合分布到了不同的节点上,而每一个子集合就称为分区(Partition)。因此也可以说,RDD是由若干个分区组成的。
RDD各个分区中的数据可以并行计算,因此分区的数量决定了并行计算的粒度。Spark会给每一个分区分配一个单独的Task任务对其进行计算,因此并行Task的数量是由分区的数量决定的。RDD分区的一个分区原则是使得分区的数量尽量等于集群中CPU核心数量。
file

04 RDD的依赖

Spark中,对RDD的每一次转化操作都会生成一个新的RDD,由于RDD的懒加载特性,新的RDD会依赖原有RDD,因此RDD之间存在类似流水线的前后依赖关系。这种依赖关系分为两种:窄依赖和宽依赖。

1、窄依赖

窄依赖是指父RDD的一个分区最多被子RDD的一个分区所用。也就是说,父RDD的分区与子RDD的分区的对应关系为一对一或多对一。例如map()、filter()、union()等操作都会产生窄依赖。
file

2、宽依赖

宽依赖是指,父RDD的一个分区被子RDD的多个分区所用。也就是说,父RDD的分区与子RDD的分区的对应关系为多对多。例如groupByKey()、reduceByKey()、sortByKey()等操作都会产生宽依赖。
file

Stage 划分

在Spark中,对每一个RDD的操作都会生成一个新的RDD,将这些RDD用带方向的直线连接起来(从父RDD连接到子RDD)会形成一个关于计算路径的有向无环图,称为DAG(Directed Acyclic Graph)。
file
Spark会根据DAG将整个计算划分为多个阶段,每个阶段称为一个Stage。每个Stage由多个Task任务并行进行计算,每个Task任务作用在一个分区上,一个Stage的总Task任务数量是由Stage中最后一个RDD的分区个数决定。

Stage的划分依据为是否有宽依赖,即是否有Shuffle。Spark调度器会从DAG图的末端向前进行递归划分,遇到Shuffle则进行划分,Shuffle之前的所有RDD组成一个Stage,整个DAG图为一个Stage。经典的单词计数执行流程的Stage划分:
file
比较复杂一点的Stage划分:
file

为说明要根据宽依赖划分Stage呢?这是因为窄依赖对优化很有利。逻辑上,每个RDD的算子都是一个fork/join操作(此join非上文的join算子,而是指同步多个并行任务的barrier): 把计算fork到每个分区,算完后join,然后fork/join下一个RDD的算子。如果子RDD的分区到 父RDD的分区是窄依赖,就可以实施优化把两个fork/join合为一个;如果连续的变换算子序列都是窄依赖,就可以把很多个 fork/join并为一个,这将极大地提升性能。Spark把这个叫做 流水线(pipeline)优化。

05 RDD的持久化

Spark中的RDD是懒加载的,只有当遇到行动算子时才会从头计算所有RDD,而且当同一个RDD被多次使用时,每次都需要重新计算一遍,这样会严重增加消耗。为了避免重复计算同一个RDD,可以将RDD进行持久化。
Spark中最重要的功能之一是可以将某个RDD中的数据保存到内存或者磁盘中,每次需要对这个RDD进行算子操作时,可以直接从内存或磁盘中取出该RDD的持久化数据,而不需要从头计算才能得到这个RDD。
例如有多个RDD,它们的依赖关系如图。若RDD3没有持久化保存,则每次对RDD3进行操作时都需要从textFile()开始计算,将文件数据转化为RDD1,再转化为RDD2,最终才得到RDD3。
可以在RDD上使用persist()或cache()方法来标记要持久化的RDD(cache()方法实际上底层调用的是persist()方法)。

file

RDD的检查点

RDD的检查点机制(Checkpoint)相当于对RDD数据进行快照,可以将经常使用的RDD快照到指定的文件系统中,最好是共享文件系统,例如HDFS。当机器发生故障导致内存或磁盘中的RDD数据丢失时可以快速从快照中对指定的RDD进行恢复,而不需要根据RDD的依赖关系从头进行计算,大大提高了计算效率。与cache()或者persist()将RDD数据存放到内存或者磁盘中的不同有以下几点:

(1)cache()或者persist()是将数据存储于机器本地的内存或磁盘,当机器故障时无法进行数据恢复,而检查点是将RDD数据存储于外部的共享文件系统(例如HDFS),共享文件系统的副本机制保证了数据的可靠性。

(2)在Spark应用程序执行结束后,cache()或者persist()存储的数据将被清空,而检查点存储的数据不会受影响,将永久存在,除非手动将其移除。因此,检查点数据可以被下一个Spark应用程序使用,而cache()或者persist()数据只能被当前Spark应用程序使用。

val sc = new SparkContext(conf);
//设置检查点数据存储路径
sc.setCheckpointDir("hdfs://centos01:8020/spark-ck")

06 共享变量

通常情况下,Spark应用程序运行的时候,Spark算子(例如map(func)或filter(func))中的函数func会被发送到远程的多个Worker节点上执行,如果一个算子中使用了某个外部变量,则该变量会拷贝到Worker节点的每一个Task任务中,各个Task任务对变量的操作相互独立。当变量所存储的数据量非常大时(例如一个大型集合)将增加网络传输及内存的开销。因此,Spark提供了两种共享变量:广播变量和累加器。

广播变量

广播变量是将一个变量通过广播的形式发送到每个Worker节点的缓存中,而不是发送到每个Task任务中,各个Task任务可以共享该变量的数据。因此广播变量是只读的。
1、默认情况下变量的传递
例如,map()算子中使用了外部变量arr:

val arr=Array(1,2,3,4,5);
val lines:RDD[String] = sc.textFile("file:///tmp/data.txt")
val result = lines.map(line =>
   (line, arr)
)

变量传递流程如图。
file
2、使用广播变量时变量的传递
例如,使用广播变量将数组arr传递给了map()算子:

val arr=Array(1,2,3,4,5);
val broadcastVar = sc.broadcast(arr)
val result = lines.map(line =>
   (line, broadcastVar) // broadcastVar为广播变量
)

传递流程如图。

file

在分布式函数中可以通过Broadcast对象的value方法访问
广播变量的值:

scala> val broadcastVar = sc.broadcast(Array(1, 2, 3))
scala> broadcastVar.value
res0: Array[Int] = Array(1, 2, 3)

累加器

累加器提供了将Worker节点的值聚合到Driver的功能,可以用于实现计数和求和。
例如,对一个整型数组进行求和,若不使用累加器,以下代码的输出结果不正确:

var sum=0 //在Driver中声明
val rdd=sc.makeRDD(Array(1,2,3,4,5))
rdd.foreach(x=>
  //在Executor中执行
  sum+=x
)
println(sum)//输出0

使用累加器对数组进行求和:

//声明一个累加器,默认初始值为0(只能在Driver端定义)
val myacc=sc.longAccumulator("My Accumulator")
val rdd=sc.makeRDD(Array(1,2,3,4,5))
rdd.foreach(x=>
   myacc.add(x)//向累加器中添加值
)
println(myacc.value)//输出15(只能在Driver端读取)

注意,累加器只能在Driver端定义,在Executor端更新。Executor端不能读取累加器的值,需要在Driver端使用value属性读取。

Views: 90

Kafka分布式消息系统

什么是Kafka

在Spark生态体系中,Kafka占有非常重要的位置。Kafka是一个使用Scala语言编写的基于ZooKeeper的高吞吐量低延迟的分布式发布与订阅消息系统,它可以实时处理大量消息数据以满足各种需求。比如基于Hadoop的批处理系统,低延迟的实时系统等。即便使用非常普通的硬件,Kafka每秒也可以处理数百万条消息,其延迟最低只有几毫秒。
在实际开发中,Kafka常常作为Spark Streaming的实时数据源,Spark Streaming从Kafka中读取实时消息进行处理,保证了数据的可靠性与实时性。二者是实时消息处理系统的重要组成部分。
那么Kafka到底是什么?简单来说,Kafka是消息中间件的一种。

举个生产者与消费者的例子:生产者生产鸡蛋,消费者消费鸡蛋。

假设消费者消费鸡蛋的时候噎住了(系统宕机了),而生产者还在生产鸡蛋,那么新生产的鸡蛋就丢失了;

再比如,生产者1秒钟生产100个鸡蛋(大交易量的情况),而消费者1秒钟只能消费50个鸡蛋,那过不了多长时间,消费者就吃不消了(消息堵塞,最终导致系统超时),导致鸡蛋又丢失了。

这个时候我们放个篮子在生产者与消费者中间,生产者生产出来的鸡蛋都放到篮子里,消费者去篮子里拿鸡蛋,这样鸡蛋就不会丢失了,这个篮子就相当于“Kafka”;鸡蛋则相当于Kafka中的消息(Message);篮子相当于存放消息的消息队列,也就是Kafka集群;

当篮子满了,鸡蛋放不下了,这时再加几个篮子,就是Kafka集群扩容。

Kafka中的一些基本概念:

消息(Message)

Kafka的数据单元被称为消息, 也称为事件(Event)。可以把消息看成是数据库里的一行数据或一条记录。为了提高效率,消息可以分组传输,每一组消息就是一个批次,分成批次传输可以减少网络开销。

服务器节点(Broker)

Kafka集群包含一个或多个服务器节点,一个独立的服务器节点被称为Broker。

主题(Topic)

每条发布到Kafka集群的消息都有一个类别,这个类别被称为主题。在物理上,不同主题的消息分开存储;在逻辑上,一个主题的消息虽然保存于一个或多个Broker上,但用户只需指定消息的主题即可生产或消费消息而不必关心消息存于何处。

分区(Partition)

为了使Kafka的吞吐率可以水平扩展,物理上把主题分成一个或多个分区。创建主题时可指定分区数量。

生产者(Producer)

负责发布消息到Kafka的Broker,实际上属于Broker的一种客户端。生产者负责选择哪些消息应该分配到哪个主题内的哪个分区。默认生产者会把消息均匀的分布到特定主题的所有分区上,但在某些情况下,生产者会将消息直接写到指定的分区。

消费者(Consumer)

从Kafka的Broker上读取消息的客户端。读取消息时需要指定读取的主题,通常消费者会订阅一个或多个主题,并按照消息生成的顺序读取他们。

不同的主题好比不同的高速公路,分区好比某条高速公路上的车道,消息就是车道上运行的车辆。如果车流量大,则拓宽车道,反之,则减少车道;而消费者就好比高速公路上的收费站,开放的收费站越多,则车辆通过速度越快。

Kafka架构

Kafka的消息传递流程如图所示。生产者将消息发送给Kafka集群,同时Kafka集群将消息转发给消费者。
一个典型的Kafka集群中包含若干生产者(数据可以是Web前端产生的页面内容或者服务器日志等)、若干Broker、若干消费者(可以是Hadoop集群、实时监控程序、数据仓库或其它服务)以及一个ZooKeeper集群。ZooKeeper用于管理和协调Broker。当Kafka系统中新增了Broker或者某个Broker故障失效时,ZooKeeper将通知生产者和消费者。生产者和消费者据此开始与其它Broker协调工作。生产者使用Push模式将消息发送到Broker,而消费者使用Pull模式从Broker订阅并消费消息。

file

file

主题与分区

Kafka通过主题对消息进行分类,一个主题可以分为多个分区,且每个分区可以存储于不同的Broker上,也就是说,一个主题可以横跨多个服务器。

对主题进行分区的好处是:允许主题消息规模超出一台服务器的文件大小上限。因为一个主题可以有多个分区,且可以存储在不同的服务器上,当一个分区的文件大小超出了所在服务器的文件大小上限时,可以动态添加其它分区,因此可以处理无限量的数据。

file

Kafka会为每个主题维护一个分区日志,记录各个分区的消息存放情况。消息以追加的方式写入到每个分区的尾部,然后以先入先出的顺序进行读取。由于一个主题包含多个分区,所以无法在整个主题范围内保证消息的顺序,但可以保证单个分区内消息的顺序。

当一条消息被发送到Broker时,会根据分区规则被存储到某个分区里。如果分区规则设置的合理,所有消息将被均匀的分配到不同的分区里,这样就实现了水平扩展。如果一个主题的消息都存放到一个文件中,则该文件所在的Broker的I/O将成为主题的性能瓶颈,而分区正好解决了这个问题。

分区中的每个记录都被分配了一个偏移量(offset),偏移量是一个连续递增的整数值,它唯一标识分区中的某个记录。而消费者只需保存该偏移量即可,当消费者客户端向Broker发起消息请求时需要携带偏移量。例如,消费者向Broker请求主题test的分区0中的偏移量从20开始的所有消息以及主题test的分区1中的偏移量从35开始的所有消。当消费者读取消息后,偏移量会线性递增。当然,消费者也可以按照任意顺序消费消息,比如读取已经消费过的历史消息(将偏移量重置到之前版本)。此外,消费者还可以指定从某个分区中一次最多返回多少条数据,防止一次返回数据太多而耗尽客户端的内存。

file

分区副本

在Kafka集群中,为了提高数据的可靠性,同一个分区可以复制多个副本分配到不同的Broker,这种方式类似于HDFS中的副本机制。如果其中一个Broker宕机,其它Broker可以接替宕机的Broker,不过生产者和消费者需要重新连接到新的Broker。

Kafka每个分区的副本都被分为两种类型:领导者副本和跟随者副本。领导者副本只有一个,其余的都是跟随者副本。所有生产者和消费者都向领导者副本发起请求,进行消息的写入与读取,而跟随者副本并不处理客户端的请求,它唯一的任务是从领导者副本复制消息,以保持与领导者副本数据及状态的一致。
如果领导者副本发生崩溃,会从其余的跟随者副本中选出一个作为新的领导者副本。

file

file

消费者组

消费者组(Consumer Group)实际上就是一组消费者的集合。每个消费者属于一个特定的消费者组(可为每个消费者指定组名称,消费者通过组名称对自己进行标识,若不指定组名称则属于默认的组)。
传统消息处理有两种模式:队列模式和发布订阅模式。队列模式是指消费者可以从一台服务器读取消息,并且每个消息只被其中一个消费者消费;发布订阅模式是指消息通过广播方式发送给所有消费者。而Kafka提供了消费者组模式,能够同时具备这两种(队列和发布订阅)模式的特点。
Kafka规定,同一消费者组内不允许多个消费者消费同一分区的消息;而不同的消费者组,可以同时消费同一分区的消息。也就是说,分区与同一个消费者组中的消费者的对应关系是多对一而不允许一对多。举个例子,如果同一个应用有100台机器,这100台机器属于同一个消费者组,则同一条消息在100台机器中只有一台能得到。如果另一个应用也需要同时消费同一个主题的消息,则需要新建一个消费者组并消费同一个主题的消息。我们已经知道,消息存储于分区中,消费者组与分区的关系如图。

file

Kafka集群环境搭建

Kafka依赖ZooKeeper集群,搭建Kafka集群之前,需要先搭建好ZooKeeper集群。ZooKeeper集群的搭建步骤此处不做过多讲解。本例依然使用三台服务器在CentOS7上搭建Kafka集群,三台服务器的主机和IP地址分别为:

centos01 192.168.170.133
centos02 192.168.170.134
centos03 192.168.170.135

1、下载解压Kafka

从Apache官网http://kafka.apache.org下载Kafka的稳定版本kafka_2.11-2.0.0.tgz。
然后将Kafka安装包上传到centos01节点的/opt/softwares目录,并解压到目录/opt/modules下:

$ tar -zxvf kafka_2.11-2.0.0.tgz -C /opt/modules/

2、修改配置文件

修改Kafka安装目录下的config/server.properties文件。在分布式环境中建议至少修改以下配置项:

broker.id=1
num.partitions=2
default.replication.factor=2
listeners=PLAINTEXT://centos01:9092
log.dirs=/opt/modules/kafka_2.11-2.0.0/kafka-logs
zookeeper.connect=centos01:2181,centos02:2181,centos03:2181

3、发送安装文件到其他节点
将centos01节点配置好的Kafka安装文件复制到centos02和centos03节点:

scp -r kafka_2.11-2.0.0/ hadoop@centos02:/opt/modules/
scp -r kafka_2.11-2.0.0/ hadoop@centos03:/opt/modules/

复制完成后,修改centos02节点的Kafka安装目录下的config/server.properties文件,修改内容如下:

broker.id=2
listeners=PLAINTEXT://centos02:9092

同理,修改centos03节点的Kafka安装目录下的config/server.properties文件,修改内容如下:

broker.id=3
listeners=PLAINTEXT://centos03:9092

4、启动ZooKeeper集群
分别在三个节点上执行以下命令,启动ZooKeeper集群(需进入ZooKeeper安装目录):

bin/zkServer.sh start

5、启动Kafka集群
分别在三个节点上执行以下命令,启动Kafka集群(需进入Kafka安装目录):

bin/kafka-server-start.sh -daemon config/server.properties

集群启动后,分别在各个节点上执行jps命令,查看启动的Java进程,若能输出如下进程信息,说明启动成功。

2848 Jps
2518 QuorumPeerMain
2795 Kafka

查看Kafka安装目录下的日志文件logs/server.log,确保运行稳定,没有抛出异常。至此,Kafka集群搭建完成。

Kafka命令行操作

1、创建主题
创建主题可以使用Kafka提供的命令工具kafka-topics.sh,此处我们创建一个名为topictest的主题,分区数为2,每个分区的副本数为2,命令如下(在Kafka集群的任意节点执行即可):

$ bin/kafka-topics.sh \
--create \
--zookeeper centos01:2181,centos02:2181,centos03:2181 \
--replication-factor 2 \
--partitions 2 \
--topic topictest

说明

  • --create:指定命令的动作是创建主题,使用该命令必须指定--topic参数。
  • --topic:所创建的主题名称。
  • --partitions:所创建主题的分区数。
  • --zookeeper:指定ZooKeeper集群的访问地址,这种方式已废弃
    • 建议改成--bootstrap-server的方式指定broker所在地址和端口,如有多个用逗号分隔。
  • --replication-factor:所创建主题的分区副本数,其值必须小于等于Kafka的节点数。
    • 如果只有一个节点,可以不指定

命令执行完毕后,若输出以下结果则表明创建主题成功:

Created Topic "topictest".

2、查询主题
创建主题成功后,可以执行以下命令,查看当前Kafka集群中存在的所有主题:

$ bin/kafka-topics.sh \
--list \
--zookeeper centos01:2181

也可以使用--describe参数查询某一个主题的详细信息。例如,查询主题topictest的详细信息,命令如下:

$ bin/kafka-topics.sh \
--describe \
--zookeeper centos01:2181 \
--topic topictest

输出结果如下:

Topic:topictest PartitionCount:2    ReplicationFactor:2 Configs:
Topic: topictest    Partition: 0    Leader: 2   Replicas: 2,3   Isr: 2,3
Topic: topictest    Partition: 1    Leader: 3   Replicas: 3,1   Isr: 3,1

可以看到,该主题有2个分区,每个分区有2个副本。分区编号为0的副本分布在broker.id为2和3的Broker上,其中broker.id为2上的副本为领导者副本;分区编号为1的副本分布在broker.id为1和3的Broker上,其中broker.id为3上的副本为领导者副本。

3、创建生产者
Kafka生产者作为消息生产角色,可以使用Kafka自带的命令工具创建一个最简单的生产者。例如,在主题topictest上创建一个生产者,命令如下:

$ bin/kafka-console-producer.sh \
--broker-list centos01:9092,centos02:9092,centos03:9092 \
--topic topictest

说明:

  • --broker-list:指定Kafka Broker的访问地址,只要能访问到其中一个即可连接成功,若想写多个则用逗号隔开。建议将所有的Broker都写上,如果只写其中一个,如果该Broker失效,连接将失败。注意此处的Broker访问端口为9092,Broker通过该端口接收生产者和消费者的请求,该端口在安装Kafka时已经指定。
  • --topic:指定生产者发送消息的主题名称。
    创建完成后,控制台进入等待键盘输入消息的状态。

4、创建消费者
新开启一个SSH连接窗口(可连接Kafka集群中的任何一个节点),在主题topictest上创建一个消费者,命令如下:

$ bin/kafka-console-consumer.sh \
--bootstrap-server centos01:9092,centos02:9092,centos03:9092 \
--topic topictest

上述代码中,参数--bootstrap-server用于指定Kafka Broker访问地址。
消费者创建完成后,等待接收生产者的消息。此时在生产者控制台输入消息“hello kafka”后按回车(可以将文件或者标准输入的消息发送到Kafka集群中,默认一行作为一个消息)即可将消息发送到Kafka集群。

file

在消费者控制台,则可以看到输出相同的消息“hello kafka”:

file

谢谢!

Views: 268

初识 Spark (含环境搭建)

大数据开发总体架构

file

什么是Spark

Apache Spark是一个快速通用的集群计算系统,是一种与Hadoop相似的开源集群计算环境,但是Spark在某些工作负载方面表现得更加优越。它提供了Java、Scala、Python和R的高级API,以及一个支持通用的执行图计算的优化引擎。它还支持一组丰富的高级工具,包括使用SQL进行结构化数据处理的Spark SQL、用于机器学习的MLlib、用于图处理的GraphX,以及用于实时流处理的Spark Streaming。

Spark的主要特点:

1、快速
与MapReduce相比,Spark可以支持包括Map和Reduce在内的更多操作,这些操作相互连接形成一个有向无环图(Directed Acyclic Graph,简称DAG),各个操作的中间数据则会被保存在内存中。因此处理速度比MapReduce更加快。Spark通过使用先进的DAG调度器、查询优化器和物理执行引擎,从而能够高性能的实现批处理和流数据处理。

file

2、易用
Spark可以使用Java、Scala、Python、R和SQL快速编写应用程序。
Spark提供了超过80个高级算子(关于算子,在第3章将详细讲解),使用这些算子可以轻松构建并行应用程序,并且可以从Scala、Python、R和SQL的Shell中交互式地使用它们。

3、通用
Spark拥有一系列库,包括SQL和DataFrame、用于机器学习的MLlib、用于图计算的GraphX、用于实时计算的Spark Streaming。可以在同一个应用程序中无缝地组合这些库。

4、到处运行
Spark可以使用独立集群模式运行(使用自带的独立资源调度器,称为Standalone模式),也可以运行在Amazon EC2、Hadoop YARN、Mesos(Apache下的一个开源分布式资源管理框架)、Kubernetes之上,并且可以访问HDFS、Cassandra、HBase、Hive等数百个数据源中的数据。

Spark主要组件

Spark是由多个组件构成的软件栈,Spark 的核心(Spark Core)是一个对由很多计算任务组成的、运行在多个工作机器或者一个计算集群上的应用进行调度、分发以及监控的计算引擎。

file

Spark运行时架构

Spark有多种运行模式,可以运行在一台机器上,称为本地(单机)模式;也可以以YARN或Mesos作为底层资源调度系统以分布式的方式在集群中运行,称为Spark On YARN模式;还可以使用Spark自带的资源调度系统,称为Spark Standalone模式。
本地模式通过多线程模拟分布式计算,通常用于对应用程序的简单测试。本地模式在提交应用程序后,将会在本地生成一个名为“SparkSubmit”的进程,该进程既负责程序的提交又负责任务的分配、执行和监控等。

YARN集群架构

在学习Spark集群架构之前,先需要了解YARN集群的架构。YARN集群总体上是经典的主/从(Master/Slave)架构,主要由ResourceManager、NodeManager、ApplicationMaster和Container等几个组件构成。

file

YARN集群中应用程序的执行流程:

file

Spark Standalone架构

Spark Standalone模式为经典的Master/Slave架构,资源调度是Spark自己实现的。在Standalone模式中,根据应用程序提交的方式不同,Driver(主控进程)在集群中的位置也有所不同。应用程序的提交方式主要有两种:client和cluster,默认是client。

当提交方式为client时,运行架构:

file

当提交方式为cluster时,运行架构:

file

Spark On YARN架构

Spark On YARN模式,遵循YARN的官方规范,YARN只负责资源的管理和调度,运行哪种应用程序由用户自己实现,因此可能在YARN上同时运行MapReduce程序和Spark程序,YARN很好的对每一个程序实现了资源的隔离。这使得Spark与MapReduce可以运行于同一个集群中,共享集群存储资源与计算资源。Spark On YARN模式与Standalone模式一样,也分为client和cluster两种提交方式。

client提交方式架构:

file

cluster提交方式架构:

file

Spark Standalone 集群搭建

Spark Standalone模式的搭建需要在集群的每个节点都安装Spark,集群角色分配如表:

file

1、下载解压安装包
访问Spark官网http://spark.apache.org/downloads.html下载预编译的Spark安装包,选择Spark版本为2.4.0,包类型为“Pre-built for Apache Hadoop 2.7 and later”(Hadoop2.7及之后版本的预编译版本)。

$ tar -zxvf spark-2.4.0-bin-hadoop2.7.tgz -C /opt/modules/

2、修改配置文件
修改slaves文件:

$ cp slaves.template slaves
$ vi slaves

改为以下内容:

centos02
Centos03

修改spark-env.sh文件:

$ cp spark-env.sh.template spark-env.sh
$ vi spark-env.sh

改为以下内容:

export JAVA_HOME=/opt/modules/jdk1.8.0_144
export SPARK_MASTER_HOST=centos01
export SPARK_MASTER_PORT=7077

启动默认的log4j日志配置

cp log4j.properties.template log4j.properties

3、拷贝安装文件到其他节点

$ scp -r /opt/modules/spark-2.4.0-bin-hadoop2.7/ hadoop@centos02:/opt/modules/
$ scp -r /opt/modules/spark-2.4.0-bin-hadoop2.7/ hadoop@centos03:/opt/modules/

4、启动Spark集群
在主节点(centos01)执行,启动Spark集群:

$ sbin/start-all.sh

使用jps查看启动进程,三个节点的进程分别为:Master、Worker、Worker说明启动成功。
访问网址http://centos01:8080,查看Spark的WebUI:

file

Spark提供了一个客户端应用程序提交工具spark-submit,使用该工具可以将编写好的Spark应用程序提交到Spark集群:

$ bin/spark-submit [options] <app jar> [app options]

说明

  • [options]:表示传递给spark-submit的控制参数;
  • <app jar>:表示提交的程序jar包(或Python脚本文件)所在位置;
  • [app options]:表示jar程序需要传递的参数,例如main()方法中需要传递的参数。

例如:

本地模式,使用2个cpu核心:
本地模式不会提交给Spark Master,因此在Spark Master WebUI http://centos01:8080/ 看不到提交的任务信息。

$ bin/spark-submit --master local[2] --deploy-mode client --class org.apache.spark.examples.SparkPi  ./examples/jars/spark-examples_2.11-2.4.8.jar

在Standalone模式(Spark集群使用Spark自带的资源协调服务)下,将Spark自带的求圆周率的程序提交到集群:

$ bin/spark-submit \
--master spark://centos01:7077 \
--class org.apache.spark.examples.SparkPi \
./examples/jars/spark-examples_2.11-2.4.0.jar

说明:
--master参数指定了Master节点的连接地址。该参数根据不同的Spark集群模式,其取值也有所不同:

在输出的日志中间可以找到Pi的估算结果:

Pi is roughly 3.138195690978455

另外在Spark任务运行时,从Spark Master UI界面可以在Running Applications选项下查看当前运行的Spark任务的状态,任务执行完毕后, 在Completed Applications选项中可以找到执行完成的Spark任务,点击Application ID可以查看任务详情和日志信息,

取值 描述
spark://host:port Standalone模式下的Master节点的连接地址,默认端口为7077
yarn 连接到YARN集群。若YARN中没有指定ResourceManager的启动地址,则需要在ResourceManager所在的节点上进行应用程序的提交,否则将因找不到ResourceManager而提交失败
local 运行本地模式,使用1个CPU核心
local[N] 运行本地模式,使用N个CPU核心。例如,local[2]表示使用2个CPU核心运行程序
local[*] 运行本地模式,尽可能使用最多的CPU核心

spark-submit还提供了一些控制资源使用和运行时环境的参数:

参数 描述
--master Master节点的连接地址。取值为spark://host:port、mesos://host:port、yarn、k8s://https://host:port或local (默认为 local[*])
--deploy-mode 提交方式。取值为“client”或“cluster”。“client”表示在本地客户端启动Driver程序,“cluster”表示在集群内部的工作节点上启动Driver程序。默认为“client”
--class 应用程序的主类(Java或Scala程序)
--name 应用程序名称,会在Spark Web UI中显示
--jars 应用依赖的第三方的jar包列表,以逗号分隔
--files 需要放到应用工作目录中的文件列表,以逗号分隔。此参数一般用来放需要分发到各节点的数据文件
--conf 设置任意的SparkConf配置属性。格式为“属性名=属性值”
--properties-file 加载外部包含键值对的属性文件。如果不指定,默认将读取Spark安装目录下的conf/spark-defaults.conf文件中的配置
--driver-memory Driver进程使用的内存量。例如“512M”或“1G”,单位不区分大小写。默认为1024M
--executor-memory 每个Executor进程所使用的内存量。例如“512M”或“1G”,单位不区分大小写。默认为1G
--driver-cores 每个Executor进程所使用的内存量。例如“512M”或“1G”,单位不区分大小写。默认为1G
--executor-cores 每个Executor进程所使用的CPU核心数,默认为1
--num-executors Executor进程数量,默认为2。如果开启动态分配,则初始Executor的数量至少是此参数配置的数量。需要注意的是,此参数仅在Spark On YARN模式中使用

例如,在Standalone模式下,将Spark自带的求圆周率的程序提交到集群,并且设置Driver进程使用内存为512M,每个Executor进程使用内存为1G,每个Executor进程所使用的CPU核心数为1,提交方式为cluster(即Driver进程运行在集群的工作节点中),执行命令:

$ bin/spark-submit \
--master spark://centos01:7077 \
--deploy-mode cluster \
--class org.apache.spark.examples.SparkPi \
--driver-memory 512m \
--executor-memory 1g \
--executor-cores 1 \
./examples/jars/spark-examples_2.11-2.4.0.jar

说明:

查看UI界面可以得知每个Worker的CPU核心数,注意设置的executor-cores不能超过这个数量。

Spark带有交互式的Shell,可在Spark Shell中直接编写Spark任务,然后提交到集群与分布式数据进行交互,并且可以立即查看输出结果。

Spark Standalone模式启动Spark Shell终端:

$ bin/spark-shell --master spark://centos01:7077

Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
Spark context Web UI available at http://hadoop100:4040
Spark context available as 'sc' (master = spark://centos01:7077, app id = app-20210727162737-0001).
Spark session available as 'spark'.
Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 2.4.0:
      /_/

Using Scala version 2.11.12 (Java HotSpot(TM) 64-Bit Server VM, Java 1.8.0_281)
Type in expressions to have them evaluated.
Type :help for more information.

scala>

启动完成后,访问Spark Master WebUI http://centos01:8080/ 查看运行的Spark应用程序:

file

退出Spark Shell:(注意命令前面以冒号:开头,可以简写为:q)

scala>:quit

Spark On YARN 集群模式搭建

Spark On YARN模式下Spark Shell的启动与Standalone模式所不同的是:--master的参数值为yarn。例如以下启动命令:

$ bin/spark-shell --master yarn

如果之前没有配置 Spark On YARN 集群模式的环境的话,这一步铁定会遇到异常的,当我们解决这些异常之后, Spark On YARN 集群模式也就自然搭建完成了。

Spark On YARN模式启动Shell出现的问题

Unable to load native-hadoop library for your platform

原因 这个只是警告,而不是错误,提示缺少对Hadoop的lib的引用。在环境变量里面进行设置即可。如果没有配置的话则会使用内建的Java类来实现,导致执行效率上有一定影响。

解决方法

编辑 /etc/profile,添加:

export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native/:$LD_LIBRARY_PATH

使环境变量生效

source /etc/profile

也可以在需要执行的脚本中前面加上

#!/bin/bash
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native/

either HADOOP_CONF_DIR or YARN_CONF_DIR must be set

解决方法
编辑spark-env.sh, 增加

export HADOOP_CONF_DIR=$HADOOP_HOME/etc/hadoop

Neither spark.yarn.jars nor spark.yarn.archive is set

$ hdfs dfs -mkdir -p /user/hadoop/spark/jars
$ hdfs dfs -put /opt/pkg/spark/jars/* /user/hadoop/spark/jars
$ hdfs dfs -chmod -R 755 /user/hadoop/spark/jars

spark-defaults.conf中写入 (注意后面的/*别漏加)

spark.yarn.jars  hdfs://centos01:8020/user/hadoop/spark/jars/*

说明:
centos01:8020 对应HDFS的NameNode的主机名和端口号(对于hadoop2.x的NameNode端口默认是9000,对于hadoop3.x的NameNode端口默认是8020)

最后重新启动Spark,然后运行Spark Shell:

$ bin/spark-shell --master yarn
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
Spark context Web UI available at http://centos01:4040
Spark context available as 'sc' (master = yarn, app id = application_1627371733487_0005).
Spark session available as 'spark'.
Welcome to
      ____              __
     / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /___/ .__/\_,_/_/ /_/\_\   version 2.4.8
      /_/

Using Scala version 2.11.12 (Java HotSpot(TM) 64-Bit Server VM, Java 1.8.0_281)
Type in expressions to have them evaluated.
Type :help for more information.

这就是正常运行的样子。

Container Killed on request.

出现这个错误说明Yarn给Spark任务的内存分配过小,Yarn最终直接将Container的进程杀掉了。

解决方法:
在Hadoop的配置文件yarn-site.xml中加入以下内容即可:

    <!--关闭物理内存检查-->
    <property>
        <name>yarn.nodemanager.pmem-check-enabled</name>
        <value>false</value>
    </property>
    <!--关闭虚拟内存检查-->
    <property>
        <name>yarn.nodemanager.vmem-check-enabled</name>
        <value>false</value>
    </property>

修改完毕后,将此文件分发到所有Yarn节点,重启Yarn集群。

Spark 集群 Master 的 HA 环境搭建

默认情况下在Spark standalone集群中进行计算时,由于是RDD的计算模型,所以可以认为worker 已经是有HA特性的了,但是负责资源调度的Master节点有可能出现单点故障。所以为了保证环境的稳定,还是需要配置HA功能。

官方文档中提供了两种HA的机制

  • 基于Zookeeper:利用ZooKeeper来提供主节点选举和集群状态的存储,可以在集群中运行多个连接到同一个Zookeeper集群的Master节点。其中一个将被选为“Leador”和其他的节点将会是Standby模式。如果当前的Leador挂了,会通过选举产生一个新的Leador。
    • 请注意,对于新运行的应用, 延迟可能会有1-2分钟,对于已经运行的应用不会受到影响。
  • 基于文件系统:将状态数据存在目录文件中,当主节点挂掉后,通过重启来解决问题。由于stop-master.sh方式停止Master节点是,不会将对应的数据删除,所以时间长了,可能会影响到启动速度

为了解决Spark的Master单点故障问题,这里基于Zookeeper实现两个Master的主备切换。

首先停止集群,修改spark-env.sh,删除SPARK_MASTER_IP属性配置(如果有),并添加以下配置

export SPARK_DAEMON_JAVA_OPTS="-Dspark.deploy.recoveryMode=ZOOKEEPER -Dspark.deploy.zookeeper.url=centos01:2181,centos02:2181,centos03:2181 -Dspark.deploy.zookeeper.dir=/spark"

然后同步spark-env.sh文件到其他节点

$ scp $SPARK_HOME/conf/spark-env.sh hadoop@centos01:$SPARK_HOME/conf/
spark-env.sh
$ scp $SPARK_HOME/conf/spark-env.sh hadoop@centos02:$SPARK_HOME/conf/
spark-env.sh

在主Master节点启动集群

$ sbin/start-all.sh

在备Master节点启动Master

$ sbin/start-master.sh

异常解决

一些可能会遇到的异常如下:

配置HA时备用Master启动失败

节点centos02启动master后,使用jps命令发现没有master进程,查看logs下面的master日志发现如下报错:
Service 'sparkMaster' failed after 16 retries

解决办法:

在备用Master所在节点的spark-env.sh中配置SPARK_LOCAL_IP属性,对应的值就是此节点的IP或者主机名(Cluster模式下也可以配置localhost或者127.0.0.1, 但是不建议)

# Options read by executors and drivers running inside the cluster
# - SPARK_LOCAL_IP, to set the IP address Spark binds to on this node
export SPARK_LOCAL_IP=centos02

spark master web ui 默认端口 8080 被占用

spark master web ui 默认端口为8080,当系统有其它程序(hadoop3.x版本的集群中的有的节点会启动jetty,用的就是8080端口)也在使用该接口时,启动master时就会报错,为了避免端口冲突,我们也可以自行设置端口号,修改方法:

方法一(推荐)

spark-env.sh中配置SPARK_MASTER_WEBUI_PORT属性,对应的值就是端口号(注意不要和已经开放的端口冲突):

# - SPARK_MASTER_WEBUI_PORT, to use non-default ports for the master
export SPARK_MASTER_WEBUI_PORT=8085
方法二

如果没有在spark-env.sh中指定SPARK Master Web UI端口,则会使用Master启动脚本sbin/start-master.sh中定义的默认端口号,即变量SPARK_MASTER_WEBUI_PORT属性,默认是8080,所以也可以修改这个默认端口号:

if [ "$SPARK_MASTER_WEBUI_PORT" = "" ]; then
  SPARK_MASTER_WEBUI_PORT=8080
fi

现在分别打开活动Master和备用Master的UI界面,可以发现一个状态是Active,另一个状态为Master,说明当前的HA配置成功。

centos01(活动 Master 节点)

URL: spark://centos01:7077
Alive Workers: 2
Cores in use: 2 Total, 0 Used
Memory in use: 4.5 GB Total, 0.0 B Used
Applications: 0 Running, 0 Completed
Drivers: 0 Running, 0 Completed
Status: ALIVE

centos02(备用 Master 节点)

URL: spark://centos02:7078
Alive Workers: 0
Cores in use: 0 Total, 0 Used
Memory in use: 0.0 B Total, 0.0 B Used
Applications: 0 Running, 0 Completed
Drivers: 0 Running, 0 Completed
Status: STANDBY

备注:
Master使用的端口是从7077开始,如果找不到就+1, 如此下去直至找到可用端口,默认最多尝试16次,如果还没有找到就报错。

Spark HA测试

接下来进行主备切换测试:

首先centos01关闭Master进程(也可以kill -9直接杀死Master进程)模拟活动Master故障的情况:

此时 http://centos01:8080/ 就不能访问了。

过几秒钟后centos02的Master Web UI界面http://centos02:8080/ 可以发现Master状态变为恢复状态 - RECOVERING,这个状态的持续时间很短,只有几秒钟:

http://centos02:8080/

Spark Master at spark://centos02:7078
URL: spark://centos02:7078
Alive Workers: 0
Cores in use: 0 Total, 0 Used
Memory in use: 0.0 B Total, 0.0 B Used
Applications: 0 Running, 0 Completed
Drivers: 0 Running, 0 Completed
Status: RECOVERING

然后很快就切换成活动状态 - Active,说明Spark的HA配置成功。
http://centos02:8080/

Spark Master at spark://centos02:7078
URL: spark://centos02:7078
Alive Workers: 2
Cores in use: 2 Total, 0 Used
Memory in use: 4.5 GB Total, 0.0 B Used
Applications: 0 Running, 0 Completed
Drivers: 0 Running, 0 Completed
Status: ALIVE

Spark 三种提交模式测试

下面的测试环境使用三台虚拟机,系统是CentOS7_x64、主机名分别使hadoop100、hadoop101、Hadoop102,安装了Java JDK8和hadoop3.1.4,Spark 2.4.8 (使用scala2.11编译的) built for Hadoop 2.7.3(用的hadoop的32位的库,因此有些兼容问题),其中Spark的活动master在hadoop100,备用master在hadoop101,worker节点是hadoop101和hadoop102。YARN的ResourceManager节点在hadoop101。

虚拟机环境介绍完了,下面编写一些脚本,使用官方提供的估算Pi值的程序,分别测试Local(本地)、Standalone(提交给Spark集群)以及Spark on YARN(提交给YARN集群)这三种提交模式,除了Local提交模式之外,根据Driver的的部署方式还分为Client模式部署以及Cluster模式部署两种部署方法。

本地客户端提交:
仅使用当前节点计算,默认使用一个CPU核心。
由于客户端和Driver在一个进程,结果直接显示在控制台。

Local[2]表示使用2个核心, Local[*]表示使用尽可能多的核心

#!/bin/bash
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master local \
--name pi-local \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

Standalone-Client部署模式:
提交任务使用spark自带的资源管理机制
计算任务交给Spark配置的工作节点(工作节点配置在slaves中).
Drive和客户端在一个进程,直接在控制台输出结果

#!/bin/bash
# spark预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master spark://hadoop100:7077 \
--total-executor-cores 2 \
--name pi-cluster-client \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

提交Python编写的Spark任务到Spark集群(虚拟机中的python版本为2.75)
这里使用的是Standalone-Client模式

#!/bin/bash
# Run a Python application on a Spark standalone-client mode
# spark2.4.8预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
  --master spark://hadoop100:7077 \
  /opt/pkg/spark/examples/src/main/python/pi.py \
  100

Standalone-Cluster提交,使用2个CPU核心
Standalone提交任务使用spark集群自带的资源管理机制
cluster部署模式:Master选择一个Workder节点(slaves.sh中配置)运行driver进程
注意这种提交方式要求应用程序使用的jar包和文件需要同步到所有worker节点(或放在HDFS上)
pi的计算结果可以访问Storm WebUI在driver的stdout日志中查找

export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master spark://hadoop101:6066 \
--deploy-mode cluster \
--total-executor-cores 2 \
--name pi-cluster-cluster \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

带监督模式的Standalone-Cluster部署

Cluster部署有个好处就是可以开启监督模式(supervise)
开启监督运行模式后当任务失败后可以自动重启

#!/bin/bath
# spark2.4.8预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master spark://hadoop100:6066 \
--deploy-mode cluster \
# --supervise 失败时自动重启driver
--supervise \
--total-executor-cores 2 \
--name pi-cluster-cluster \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

spark以client方式提交时,port应该设置为7077;以cluster方式提交时,port设置为6066,因为这种方式提交时,是以rest api方式提交application。

YARN-Client部署模式:
客户端运行driver,任务提交到YARN集群进行调度. 使用YARN集群的工作节点进行计算和(和slaves的配置无关)。可以在YARN的可视化界面(默认端口8088)查看任务执行情况,得到pi的计算结果再发送给Driver在控制台显示

#!/bin/bash
# spark预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master yarn \
--deploy-mode client \
--total-executor-cores 2 \
--name pi-yarn-yarn \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

YARN-cluster部署模式
客户端将任务交给YARN集群进行调度
使用YARN集群中的工作节点和spark中的slaves.sh配置无关
driver运行在YARN中的某个NM节点上
可以在YARN的可视化界面(默认端口8088)查看任务执行情况
结果在Application的stdout日志中查看(需提前启动JobHsotry Server)
或者使用yarn logs --applicationId <applicationId>查看

#!/bin/bash
# YARN-cluster部署模式:
# spark预编译的hadoop版本是32位的和虚拟机里的不兼容
# 因此需要指定hadoop的本地库
export LD_LIBRARY_PATH=$HADOOP_HOME/lib/native
/opt/pkg/spark/bin/spark-submit \
--class org.apache.spark.examples.SparkPi \
--master yarn \
--deploy-mode cluster \
--total-executor-cores 2 \
--name pi-yarn-yarn \
/opt/pkg/spark/examples/jars/spark-examples_2.11-2.4.8.jar \
100

Spark Shell的使用

Spark Shell的使用也有三种模式

  1. 本地(单机)模式下启动Spark Shell

    即不加--master参数,直接使用bin/spark-shell命令启动 Spark Shell,在本地模式下,所有操作任务只是在本地,也就是当前节点运行,而不会分发到整个集群。

  2. Spark Standalone模式下启动Spark Shell
    使用任意Spark节点进入Spark安装目录,执行以下命令,自动Spark Shell客户端:

$ bin/spark-shell --master spark://hadoop100:7077

说明:
--master指定Master节点的访问地址,centos01为Master所在节点主机名,7077为Master默认端口。

在Spark Shell启动过程日志中可以看出,有个Spark的上下文变量叫做sc,这个变量可以在Spark Shell中直接使用,它也是Spark应用程序的入口,负责于Spark集群进行交互。
启动完成后可以在 http://centos01:8080/ 查看运行的Spark应用程序。

  1. Spark on YARN 模式下启动Spark Shell

前提,需要启动Hadoop的HDFS文件系统和Yarn的相关进程,并完成Spark on YARN 模式的Spark环境搭建。启动命令如下:

$ bin/spark-shell --master yarn

Views: 171

Scala语言基础

什么是Scala

Scala是一种将面向对象和函数式编程结合在一起的高级语言,旨在以简洁、优雅和类型安全的方式表达通用编程模式。Scala功能强大,不仅可以编写简单脚本,还可以构建大型系统。

Scala运行于Java平台,Scala程序会通过JVM被编译成class字节码文件,然后在操作系统上运行。其运行时候的性能通常与Java程序不分上下,并且Scala代码可以调用Java方法、继承Java类、实现Java接口等,几乎所有Scala代码都大量使用了Java类库。

由于Spark主要是由Scala语言编写的,为了后续更好的学习Spark以及使用Scala编写Spark应用程序,需要首先学习使用Scala语言。

安装Scala

由于Scala运行于Java平台,因此安装Scala之前需要确保系统安装了JDK。此处使用的Scala版本为2.12.7,要求JDK版本为1.8。

Windows安装Scala

1、下载Scala
到Scala官网https://www.scala-lang.org/download/下载Windows安装包scala-2.12.7.msi

2、配置环境变量
变量名:SCALA_HOME
变量值:C:\Program Files (x86)\scala
变量名:Path
变量值:%SCALA_HOME%\bin

3、测试
CMD中执行scala -version命令

file

CentOS7安装Scala

1、下载Scala
到Scala官网https://www.scala-lang.org/download/下载Linux安装包scala-2.12.7.tgz

解压到指定目录:

$ tar -zxvf scala-2.12.7.tgz -C /opt/modules/

2、配置环境变量

export SCALA_HOME=/opt/modules/scala-2.12.7/
export PATH=$PATH:$SCALA_HOME/bin

3、测试
CMD中执行scala -version命令

Scala基础

最初学习Scala的时候建议在Scala命令行模式中操作,最终程序的编写可以在IDE中进行。在Windows CMD窗口中或CentOS的Shell命令中执行“scala”命令,即可进入Scala的命令行操作模式(REPL)。

变量声明

Scala中变量的声明使用关键字val和var。

声明一个val字符串变量str:

scala> val str="hello scala"
str: String = hello scala

声明变量时指定数据类型:

scala> val str:String="hello scala"
str: String = hello scala

将多个变量放在一起进行声明:

scala> val x,y="hello scala"
x: String = hello scala
y: String = hello scala

Scala变量的声明,需要注意的地方总结如下:

  • 定义变量需要初始化,否则会报错。
  • 定义变量时可以不指定数据类型,系统会根据初始化值推断变量的类型。
  • Scala中鼓励优先使用val(常量),除非确实需要对其进行修改。
  • Scala语句不需要写结束符,除非同一行代码使用多条语句时才需要使用分号隔开。

数据类型

在Scala中,所有的值都有一个类型,包括数值和函数。

file

Any是Scala类层次结构的根,也被称为超类或顶级类。Scala执行环境中的每个类都直接或间接地从该类继承。该类中定义了一些通用的方法,例如equals()、hashCode()和toString()。Any有两个直接子类:AnyVal和AnyRef。

AnyVal表示值类型。有9种预定义的值类型,它们是非空的Double、Float、Long、Int、Short、Byte、Char、Unit和Boolean。Unit是一个不包含任何信息的值类型,和Java语言中的void等同,用作不返回任何结果的方法的结果类型。Unit只有一个实例值,写成()。

AnyRef表示引用类型。所有非值类型都被定义为引用类型。Scala中的每个用户定义类型都是AnyRef的子类型。AnyRef对应于Java中的Java.lang.Object。

Nothing是所有类型的子类,在Scala类层级的最低端。Nothing没有对象,因此没有具体值,但是可以用来定义一个空类型,类似于Java中的标示性接口(如:Serializable,用来标识该类可以进行序列化)。举个例子,如果一个方法抛出异常,则异常的返回值类型就是Nothing(虽然不会返回)。

Null是所有引用类型(AnyRef)的子类,所以Null可以赋值给所有的引用类型,但不能赋值给值类型,这个和Java的语义是相同的。Null有一个唯一的单例值null。

下面的例子定义了一个类型为List[Any]的变量list,list中包括字符串、整数、字符、布尔值和函数,由于这些元素都属于对象Any的实例,因此可以将它们添加到list中。

val list: List[Any] = List(
  "a string",
  732,  //an integer
  'c',  //a character
  true, //a boolean value
  () => "an anonymous function returning a string"
)

Scala中的值类型可以进行转换,且转换是单向的。

file

例如下面的例子,允许将Long型转换为Float型,Char型转换为Int型:

val x: Long = 987654321
val y: Float = x  //9.8765434E8 (注意在这种情况下会丢失一些精度)

val face: Char = '☺'
val number: Int = face  //9786

表达式

Scala中常用的表达式主要有条件表达式和块表达式。

1、条件表达式
条件表达式主要是含有if/else的语句块:

scala> val i=1
i: Int = 1
scala> val result=if(i>0) 100 else -100
result: Int = 100

也可以在一个表达式中进行多次判断:

scala> val result=if(i>0) 100 else if(i==0) 50 else 10
result: Int = 100

2、块表达式
块表达式为包含在符号{}中的语句块:

scala> val result={
     | val a=10
     | val b=10
     | a+b
     | }
result: Int = 20

Scala中的返回值是最后一条语句的执行结果,而不需要像Java一样单独写return关键字。如果表达式中没有执行结果,则返回一个Unit对象,类似Java中的void:

scala> val result={
     | val a=10
     | }
result: Unit = ()

循环

Scala中的循环主要有for循环、while循环和do while循环三种。
1、for循环
for循环的语法:
for(变量<-集合或数组){
方法体
}

例如,循环从1到5输出变量i的值:

scala> for(i<- 1 to 5) println(i)

若不想包括5,可使用关键字until:

scala> for(i<- 1 until 5) println(i)

将字符串“hello”中的字符循环输出:

scala> val str="hello"
scala> for(i<-0 until str.length) println(str(i))

将字符串看做一个由多个字符组成的集合,简化写法:

scala> for(i<-str) println(i)

2、while循环
while循环的语法:
while(条件)
{
循环体
}

例如:

scala> var i=1
i: Int = 1

scala> while(i<5){
     |  i=i+1
     |  println(i)
     | }

3、do while循环
do while 循环与while循环类似,但是do while循环会确保至少执行一次循环。语法:

do {
   循环体
} while(条件)

例如:

scala> do{
     |    i=i+1
     |    println(i)
     | }while(i<5)

方法与函数

Scala中有方法与函数。Scala 方法是类或对象中定义的成员,而函数是一个对象,可以将函数赋值给一个变量。换句话说,方法是函数的特殊形式。
1、方法
方法的定义使用def关键字,语法:
def 方法名 (参数列表):返回类型={
方法体
}
例如,将两个数字求和然后返回,返回类型为Int:

def addNum( a:Int, b:Int ) : Int = {
      var sum = 0
      sum = a + b
      return sum
}

代码简写,去掉返回类型和return关键字:

def addNum( a:Int, b:Int ) = {
      var sum = 0
      sum = a + b
      sum
}

如果方法没有返回结果,可以将返回类型设置为Unit,类似Java中的void:

def addNum( a:Int, b:Int ) : Unit = {
      var sum = 0
      sum = a + b
      println(sum)
}

在定义方法参数时,可以为某个参数指定默认值,在方法被调用时可以不为带有默认值的参数传入实参:

def addNum( a:Int=5, b:Int ) = {
      var sum = 0
      sum = a + b
      sum
}

方法的调用,通过指定参数名称,只传入参数b:
addNum(b=2)
也可以将a,b两个参数都传入:
addNum(1,2)

2、函数
函数的定义与方法不一样,语法:
(参数列表)=>函数体
定义一个匿名函数,参数为a和b,且都是Int类型,函数体为a+b:
( a:Int, b:Int ) =>a+b
如果函数体有多行,可以将函数体放入一对{}中,并且可以通过一个变量来引用函数,变量相当于函数名称:
val f1=( a:Int, b:Int ) =>{ a+b }
对上述函数进行调用:
f1(1,2)
函数也可以没有参数:
val f2=( ) =>println("hello scala")
对上述函数进行调用:
f2()

3、方法与函数的区别
(1)方法是类的一部分,而函数是一个对象并且可以赋值给一个变量。
(2)函数可以作为参数传入到方法中。
例如,定义一个方法m1,参数f要求是一个函数,该函数有两个Int类型参数,且函数的返回类型为Int,方法体中直接调用该函数:

def m1(f: (Int, Int) => Int): Int = {
  f(2, 6)
 }

定义一个函数f1:
val f1 = (x: Int, y: Int) => x + y
调用方法m1,并传入函数f1:

val res = m1(f1)
println(res)

输出结果为8。

(3)方法可以转换为函数
当把一个方法作为参数传递给其它的方法或者函数时,系统将自动将该方法转换为函数。
例如,有一个方法m2:
def m2(x:Int,y:Int) = x+y
调用(2)中的m1方法,并将m2作为参数传入,此时系统会自动将m2方法转为函数:

val res = m1(m2)
println(res)

输出结果为8。
除了系统自动转换外,也可以手动进行转换。在方法名称后加入一个空格和一个下划线,即可将方法转换为函数:

val f2=m2 _
val res=m1(f2)
println(res)

输出结果为8。

集合

Scala集合分为可变集合和不可变集合。可变集合可以对其中的元素进行修改、添加、移除;而不可变集合,永远不会改变,但是仍然可以模拟添加、移除或更新操作。这些操作都会返回一个新的集合,原集合的内容不发生改变。

Array数组

Scala中的数组分为定长数组和变长数组,定长数组初始化后不可对数组长度进行修改,而变长数组则可以修改。
1、定长数组
定义数组的同时可以初始化数据:

val arr=Array(1,2,3)//自动推断数组类型

或者

val arr=Array[Int](1,2,3)//手动指定数据类型

也可以定义时指定数组长度,稍后对其添加数据:

val arr=new Array[Int](3)
arr(0)=1
arr(1)=2
arr(2)=3

可以使用for循环对数组进行遍历:

val arr=Array(1,2,3)
for(i<-arr){
  println(i)
}

Scala对数组提供了很多常用的方法,使用起来非常方便:

val arr=Array(1,2,3)
//求数组中所有数值的和
val arrSum=arr.sum
//求数组中的最大值
val arrMax=arr.max
//求数组中的最小值
val arrMin=arr.min
//对数组进行升序排序
val arrSorted=arr.sorted
//对数组进行降序排序
val arrReverse=arr.sorted.reverse

2、变长数组
变长数组使用类scala.collection.mutable.ArrayBuffer进行定义:

//定义一个变长Int类型数组
val arr=new ArrayBuffer[Int]()
//向其中添加三个元素
arr+=1
arr+=2
arr+=3

也可以使用-=符号对变长数组中的元素进行删减,例如,去掉数组arr中值为3的元素:

arr-=3

若数组中有多个值为3的元素,将从前向后删除第一个匹配的值。
在数组arr的下标为0的位置插入两个元素1和2:
// arr.insert(0,1,2) // 经测试scala2.13.6版本这个方法没有第三个参数,只能:

arr.insert(0,1)
arr.insert(0,2)

从数组arr的下标为1的位置开始移除两个元素:

arr.remove(1, 2)

List列表

Scala中的List分为可变List和不可变List,默认使用的List为不可变List。不可变List也可以增加元素,但实际上生成了一个新的List,原List不变。
1、不可变List
创建一个Int类型的List,名为nums:

val nums: List[Int] = List(1, 2, 3, 4)

在该List的头部追加一个元素1,生成一个新的List:

val nums2=nums.+:(1)

在该List的尾部追加一个元素5,生成一个新的List:

val nums3=nums:+5

List也支持合并操作,将两个List合并为一个新的List:

val nums1: List[Int] = List(1, 2, 3)
val nums2: List[Int] = List(4, 5, 6)
val nums3=nums1++:nums2
println(nums3)

输出结果:

List(1, 2, 3, 4, 5, 6)

Map映射

Scala中的Map也分可变的Map和不可变的Map,默认为不可变Map。
1、不可变Map
创建一个不可变Map:

val mp = Map(
   "key1" -> "value1",
   "key2" -> "value2",
   "key3" -> "value3"
)

也可以使用以下写法:

val mp = Map(
   ("key1" , "value1"),
   ("key2" , "value2"),
   ("key3" , "value3")
)

循环输出上述Map中的键值数据:

for((k,v)<-mp){
   println(k+":"+v)
}

2、可变Map
创建可变Map需要引入类scala.collection.mutable.Map,创建方式与不可变Map相同。访问Map中key1的值,代码:

val mp = Map(
   ("key1" , "value1"),
   ("key2" , "value2")
)
println(mp("key1"))

修改键key1的值为value2,代码:

mp("key1")="value2"

上述代码当key1存在时执行修改操作,若key1不存在则执行添加操作。
向Map中添加元素也可以使用+=符号:

mp+=("key3" -> "value3")

mp+=(("key3","value3"))

相对应的,从Map中删除一个元素可以使用-=符号:

mp-="key3"

Tuple元组

元组是一个可以存放不同类型对象的集合,元组中的元素不可以修改。
1、定义元组
定义一个元组t:

val t=(1,"scala",2.6)

或使用以下方式,其中Tuple3是一个元组类,代表元组的长度为3:

val t2 = new Tuple3(1,"scala",2.6)

2、访问元组
可以使用方法_1、_2、_3访问其中的元素,例如,取出元组中第一个元素:

println(t._1)

和数组、字符串的位置不同,元组的元素下标从1开始。
3、迭代元组
使用 Tuple.productIterator() 方法可以迭代输出元组的所有元素:

val t = (4,3,2,1)
t.productIterator.foreach{ i =>println("Value = " + i )}

Set哈希表

Set集合存储的对象不可重复。Set集合分为可变集合和不可变集合,默认情况下使用的是不可变集合,如果要使用可变集合,则需要引用 scala.collection.mutable.Set 包。
1、定义Set
定义一个不可变集合:

val set = Set(1,2,3)

2、元素增减
与List集合一样,对于不可变Set进行元素的增加和删除,实际上会产生一个新的Set,原来的Set并没有改变:

//定义一个不可变set集合
val set = Set(1,2,3)
//增加一个元素
val set1=set+4
//减少一个元素
val set2=set-3

3、常用方法

val site = Set("Ali", "Google", "Baidu") 
println(site.head)  // 输出第一个元素
val set2=site.tail  // 取得除了第一个元素的所有元素的集合
println(set2)
println(site.isEmpty) //查看元素是否为空

使用++运算符可以连接两个集合:

val site1 = Set("Ali", "Google", "Baidu")
val site2 = Set("Faceboook", "Taobao")
val site=site1++site2
println(site)

输出结果:

Set(Faceboook, Taobao, Google, Ali, Baidu)

val num = Set(5,8,7,20,10,66)
println(num.min)  //输出集合中的最小元素
println(num.max) //输出集合中的最大元素

类和对象

1、类的定义
对象是类的具体实例,类是抽象的,不占用内存,而对象是具体的,占用存储空间。Scala中一个最简单的类定义是使用关键字class,类名必须大写。类中的方法用关键字def定义,代码:

class User{
   private var age=20
   def count(){
      age+=1
   }
}

如果一个类不写访问修饰符,则默认访问级别为Public。这与Java是不一样的。关键字new用于创建类的实例。例如,调用上述代码中的count()方法:

new User().count()

2、单例对象
Scala中没有静态方法或静态字段,但是可以使用关键字object定义一个单例对象,单例对象中的方法相当于Java中的静态方法,可以直接使用“单例对象名.方法名”方式进行调用。单例对象除了没有构造器参数外,可以拥有类的所有特性。
例如,定义一个单例对象Person,该对象中定义了一个方法showInfo():

object Person{
  private var name="zhangsan"
  private var age=20
  def showInfo():Unit={
    println("姓名:"+name+",年龄:"+age)
  }
}

可以在任何类或对象中使用代码Person.showInfo()对方法showInfo()进行调用。

3、伴生对象
当单例对象的名称与某个类的名称一样时,该对象被称为这个类的伴生对象。类被称为该对象的伴生类。类和它的伴生对象必须定义在同一个文件中,且两者可以互相访问其私有成员。例如以下代码:

class Person() {
  private var name="zhangsan"
  def showInfo(){
    println("年龄:"+Person.age) //访问伴生对象的私有成员
  }
}
object Person{
  private var age=20
  def main(args: Array[String]): Unit = {
    var per=new Person()
    println("姓名:"+per.name) //访问伴生类的私有成员
    per.showInfo()
  }
}

4、get和set方法
Scala默认会根据类的属性的修饰符生成不同的get和set方法,生成原则:
val修饰的属性,系统会自动生成一个私有常量属性和一个公有get方法。
var修饰的属性,系统会自动生成一个私有变量和一对公有get/set方法。
private var修饰的属性,系统会自动生成一对私有get/set方法,相当于类的私有属性,只能在类的内部和伴生对象中使用。
private[this]修饰的属性,系统不会生成get/set方法,即只能在类的内部使用该属性。

在Scala中,get和set方法并非被命名为getName和setName,而是被命名为name和name_=,由于JVM不允许在方法名中出现=,因此=被翻译成$eq。

private String name = "zhangsan";
  public String name() {
    return this.name;
  }
  public void name_$eq(String x$1) {
    this.name = x$1;
  }

除了系统自动生成get和set方法外,也可以手动进行编写:

class Person {
  //声明私有变量
  private var privateName="zhangsan"
  def name=privateName //定义get方法
  def name_=(name:String): Unit ={ //定义set方法
    this.privateName=name
  }
 
}
object Test{
  def main(args: Array[String]): Unit = {
    var per:Person=new Person()
    //访问变量
    per.name=“lisi” //修改
    println(per.name) //读取
  }
}

5、构造器
Scala中的构造器分为主构造器和辅助构造器。
主构造器的参数直接放在类名之后,且将被编译为类的成员变量,其值由初始化类时进行传入:

//定义主构造器,年龄age默认为18
class Person(val name:String,var age:Int=18) {

}
object Person{
  def main(args: Array[String]): Unit = {
    //调用构造器并设置name和age字段
    var per=new Person("zhangsan",20)
    println(per.age)
    println(per.name)
    per.name="lisi"//错误,val修饰的变量不可修改
  }
}

将参数age设置为私有的,参数name设置为不可修改(val):

class Person(val name:String, private var age:Int) {
}

5、构造器
构造参数也可以不带val或var,此时默认为private[this] val:

class Person(name:String,age:Int) {
}

如果需要将整个主构造器设置为私有的,只需要添加private关键字即可:

class Person private(var name:String,var age:Int) {
}

除了可以有主构造器外,还可以有任意多个辅助构造器。辅助构造器的定义需要注意以下几项:
(1)辅助构造器的方法名称为this。
(2)每一个辅助构造器的方法体中必须首先调用其它已定义的构造器。
(3)辅助构造器的参数不能使用var或val进行修饰。

定义两个辅助构造器:

class Person {
  private var name="zhangsan"
  private var age=20
  //定义辅助构造器一
  def this(name:String){
    this()//调用主构造器
    this.name=name
  }
  //定义辅助构造器二
  def this(name:String,age:Int){
    this(name)//调用辅助构造器一
    this.age=age
  }
}

上述构造器可以使用如下三种方式进行调用:

var per1=new Person//调用无参主构造器
var per2=new Person("lisi")//调用辅助构造器一
var per3=new Person("lisi",28)//调用辅助构造器二

除此之外,主构造器还可以与辅助构造器同时使用,在这种情况下,一般辅助构造器的参数要多于主构造器:

//定义主构造器
class Person(var name:String,var age:Int) {
  private var gender=""
  //定义辅助构造器
  def this(name:String,age:Int,gender:String){
    this(name,age)//调用主构造器
    this.gender=gender
  }
}
object Person{
  def main(args: Array[String]): Unit = {
    //调用辅助构造器
    var per=new Person("zhangsan",20,"male")
    println(per.name)
    println(per.age)
    println(per.gender)
  }
}

输出结果

zhangsan
20
male

抽象类

Scala的抽象类使用关键字abstract定义,具有以下特征:
(1)抽象类不能被实例化。
(2)抽象类中可以定义抽象字段(没有初始化的字段)和抽象方法(没有被实现的方法),也可以定义被初始化的字段和被实现的方法。
(3)若某个子类继承了一个抽象类,则必须实现抽象类中的抽象字段和抽象方法。且实现的过程中可以添加override关键字也可以省略。若重写了抽象类中已经实现的方法,则必须添加override关键字。
定义一个抽象类Person:

abstract class Person {
  var name:String //抽象字段
  var age:Int
  var address:String="北京"  //普通字段
  def speak() //抽象方法
  def eat():Unit={ //普通方法
    println("吃东西")
  }
}

定义一个普通类Teacher,并继承抽象类Person,实现Person中的抽象字段和抽象方法,并重写方法eat():

//继承了抽象类Person
class Teacher extends Person{
  //实现抽象字段
  var name: String = "王丽"
  var age: Int = 28
  //实现抽象方法
  def speak(): Unit = {
    println("姓名:"+this.name)
    println("年龄:"+this.age)
    println("地址:"+this.address)//继承而来
    println("擅长讲课")
  }  
  //重写非抽象方法,必须添加override关键字
  override def eat():Unit={
    println("爱吃中餐")
  }
}

定义一个测试对象,调用Teacher类中的方法,代码:

object AppTest{
  def main(args: Array[String]): Unit = {
    val teacher=new Teacher()
    //调用方法
    teacher.speak()
    teacher.eat()
  }
}

输出结果:

姓名:王丽
年龄:28
地址:北京
擅长讲课
爱吃中餐

Trait特质

Scala特质使用关键字trait定义,类似Java中使用interface定义的接口。特质除了有Java接口的功能外,还有一些特殊的功能。定义了一个特质Pet:

//定义特质(宠物)
trait Pet {
  var name:String //抽象字段
  var age:Int
  def run //抽象方法
  def eat: Unit ={ //非抽象方法
    println("吃东西")
  }
}

定义一个普通类Cat,实现了上述特质Pet(必须实现未实现的字段和方法):

class Cat extends Pet
  var name:String="john" { //实现抽象字段
  var age:Int=3
  def run: Unit = {//实现抽象方法
    println("会跑")
  }
override def eat: Unit ={ //重写非抽象方法
    println("吃鱼")
  }
}

若需要实现多个特质,可以通过with关键字添加额外特质,但位于最左侧的特质必须使用extends关键字:

trait Animal{
}
trait Runable{
}
//类Dog实现了三个特质
class Dog extends Pet with Animal with Runable{
  //省略...
}

在类实例化的时候,可以通过with关键字混入多个特质,从而使用特质中的方法。例如,定义两个特质Runable、Flyable和一个类Bird:

//定义两个特质
trait Runable{
  def run=println("会跑")
}
trait Flyable{
  def fly=println("会飞")
}
//定义一个类
class Bird{
}

在类Bird实例化时混入特质Runable和Flyable:

val bird=new Bird() with Runable with Flyable
bird.run //输出结果“会跑”
bird.fly //输出结果“会飞”

谢谢!

Views: 83