目录
Spark共提供3种基本的数据类型,分别是Spark Core引擎对应的RDD,以及Spark SQL引擎对应的DataFrame 和 DataSet。
那么三种数据类型的定义是什么,他们有什么区别呢?
1、RDD是什么
- 分片的集合
RDD(resilient distributed dataset)即弹性分布式数据集,它由多个partition(分区)组成的集合。每个partition运行在分布式集群的不同节点上。
默认情况下,每个HDFS的分区文件(默认分区文件块大小是64M)都会创建一个RDD分区。
注:一个文件被切成N份,那么每一份就代表一个分区, 每个分区能被并行的处理。
-
函数作用在每个分区
RDD共包含两种计算方式,一种是transformations转换(如:map, filter, groupBy, join等),一个种是actions操作(如:count, collect, save等),每种计算方式对应的具体函数,详见《RDD常用操作》
函数作用在每个分区上并行执行。通过转换或者操作后,生成新的RDD,旧的RDD内容不变。
注意:Transformations转换是Lazy的,也就是说从一个RDD转换生成另一个RDD的操作不是马上执行,Spark在遇到Transformations转换时只会记录需要这样的转换,并不会去执行,需要等到有Actions操作的时候才会真正启动计算过程进行计算。写代码的时候,一定要特别注意,基本是新手必踩的坑!!!
-
RDD之间有一系列依赖关系
由于RDD每次转换计算都会生成新的RDD,所以RDD会形成类似流水线一样的前后依赖关系。
Spark会记录RDD前后依赖关系,所以当有分区的数据丢失时, Spark会通过依赖关系进行重新计算,从而计算出丢失的数据,而不是对RDD所有的分区进行重新计算。
依赖关系共包含两种:宽依赖和窄依赖。
a、窄依赖是指父RDD的每个分区只被子RDD的一个分区所使用
b、宽依赖是指父RDD的每个分区都可能被多个子RDD分区所使用
如下图,一个HDFS文件被切成3个分区,函数作用于每个分区上,具体计算过程如下:
a、进行一系列Transformations转换,如下的map、filter处理后,生成了3个新的RDD(窄依赖),经过reduceByKey处理后,生成2个新的RDD(宽依赖)
b、最后,经过一个actions操作,如下的saveAsTextFile操作处理,触发整条RDD依赖链执行,将结果保存到目标HDFS中
2、Dataset是什么
Dataset 也是弹性分布式数据集,但相比于RDD,多了shema信息(类似数据表的列名),使得RDD每一行的数据, 结构都是一样的。
Spark SQL通过schema可以清楚地知道该数据集中包含哪些列,每列的名称和类型各是什么,因此,可以通过类SQL的语法对数据进行操作(如筛选某列,列聚合等),简单易用。 详见《Dataset的常见操作》
如下图,直观地体现了Dataset和RDD的区别。左侧的JavaRDD<Person>虽然以Person类为类型参数,但Spark Core框架本身不了解Person类的内部结构。而右侧的Dataset<Person>却提供了详细的结构信息,使得Spark SQL可以清楚地知道该数据集中包含哪些列,每列的名称和类型各是什么,如图中红框框起来的部分。
因此,Dataset既可以进行一些RDD操作(如map、filter等) 又可以进行一些 SQL的操作(如果group、select等)。
3、DataFrame是什么
Dataset每一行存储的是强类型的值,例如:Dataset<String>,Dataset<Person>,因此取得每条数据某个值时,可以直接使用 类的方法获取,如person.getName() ,并且Dataset可以在编译时检查类型
而DataFrame被声明为Dataset<Row>,存储的Row是无类型的,是以列名来作处理的,因此取得每条数据某个值时,要使用row.getString(0)或col("name")这样的方式解析来取得,无法知道某个值的具体的数据类型。
DataFrame的常见操作详见 《DataFrame的常见操作》
4、使用时候怎么选
使用的时候,如果数据源是HDFS(每行存储的不是json串),推荐使用RDD
如果数据源是HDFS(每行存储json串)、HIVE表(如公司的UDW)或者 外部的数据库表(如mysql、redis)等,推荐使用 DataFrame 和 Dataset
5、三者的入口类
Spark2.0以前,使用不同的数据类型,需要采用不同的入口类,如RDD对应的入口类是SparkContext,DataFrame对应的入口类是SQLContext或者HiveContext,入口类用于指定Spark集群的各种参数,并与资源管理器做交互。
Spark2.0开始及之后,提出了一个统一的入口类SparkSession,封装了SparkContext、SQLContext和HiveContext.
a、如果要创建RDD,需要通过SparkSession入口类获取SparkContext入口类,如下图:
b、如果要创建DataFrame和DataSet,通过 SparkSession入口类即可。
6、三者的转化
三种数据类型可以互相转化,常用的三种的转化操作如下:
#1)RDD 转 DataFrame 详见 《RDD常用操作》
#2)DataFrame 转 Dataset 详见《DataFrame的常见操作》
#3)Dataset 转 RDD 详见 《Dataset的常见操作》