> For the complete documentation index, see [llms.txt](https://2327257266.gitbook.io/learn-spa/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://2327257266.gitbook.io/learn-spa/spark-sql/spark-sql-zhi-dataframe.md).

# DataFrame与DataSet

&#x20;        前面我们介绍过Spark对RDD进行封装得到的一个分布式的数据集，它与RDD最大的不同在于提供的更像是一个传统的数据库里面的表，他除了数据之外还能够知道更多的信息，比如说列名、列值和列的属性，这一点就和hive很类似了，而且他也能够支持一些复杂的数据格式。同时DataFrame可以从不同来源的数组构造，例如Hive表，结构化数据文件，外部数据库或现有RDD。

## **1、DataFrame的创建**

在Spark SQL中SparkSession是创建DataFrame和执行sql的入口，创建DataFrame有三种方式：通过Spark的数据源进行创建；从Hive Table进行查询返回；从一个存在的RDD进行转换。

### 1）从Spark数据源进行创建

spark支持创建数据源文件的格式主要有 **csv format jdbc json load option options orc parquet schema table text textFile。**&#x5728;spark-shell中为我们提供了一个SparkSession对象叫做spark，我们先通过这个对象初步了解一下Spark SQL。通过spark.read. 查看spark支持创建数据源文件的格式

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdS4jTBA3jr2Xhh6gA%2F-MFdSjK7c3H-Tmf7Ienc%2Fimage.png?alt=media\&token=3984501d-65d3-4bff-ae71-f1d194065cb2)

我们举例读取一个json文件，用DataFrame来表示：

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdT27uwARATSTGNtCp%2F-MFdUYqoITUQ0LadUFfU%2Fimage.png?alt=media\&token=5d08af6f-4789-42da-9937-dcc06050de83)

通过spark.read.json从json文件中读取数据，并返回一个DataFrame类型的对象，通过这个对象可以查看文件中的数据。

### 2）从Hive Table进行查询返回

陆续补充。。。。。。

### 3）从RDD进行转换

陆续补充。。。。。。

## 2、SQL语法

SQL语法风格指的是我们查询数据的时候使用SQL语句来查询，这种风格的查询必须要有临时视图或者全局视图来辅助。下面举例来说明：

### 1）从json文件中读取数据

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdWYAhB5fGc80pzSm1%2F-MFdWd9IeyZ7HwpiPGxO%2Fimage.png?alt=media\&token=ac0bc0f5-af3e-4f4d-8b49-6a357894fbbd)

**返回一个DataFrame对象，其中包含json文件中结构化数据。接下来我们使用DataFrame对象中的方法来创建视图。**

### 2）创建一个全局临时视图

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdYdCbN0iaKMmJDJ43%2F-MFdZbfCg2jcjz3D5r5a%2Fimage.png?alt=media\&token=3f4bdb84-b684-41e6-8880-cc2fe738af8b)

**通过createOrReplaceTempView()方法为df对象创建一个临时的视图，为SQL查询提供。**

### 3）通过SQL语句实现查询全表

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdYdCbN0iaKMmJDJ43%2F-MFdZhrUFvmlklEkiDBd%2Fimage.png?alt=media\&token=0fd10a72-9dce-457d-8b90-2eef0789315b)

上面我们通过createOrReplaceTempView()方法创建临时视图，通过spark环境对象提供的sql()方法，对临时视图进行查询，这里之所以叫视图和数据库中的视图的概念类似，我们得到的是SQL查询的结果，即只能用于查询，无法对原数据进行增加和修改的操作。

{% hint style="info" %}
注意，对于普通临时视图是在Session范围内的，如果想要应用范围内有效，可以使用全局临时视图(对应的方法是createOrReplaceGlobalTempView())。此时如果访问全局临时视图需要全路径访问，如：global\_temp.user
{% endhint %}

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFd_L_6TT-zaIVRd2Te%2F-MFd_c665pyT8qHnqrW_%2Fimage.png?alt=media\&token=8a5934bb-4157-4660-81a4-c3a49ddbda39)

我们先通过createOrReplaceGlobalTempView()方法创建一个全局临时视图，在通过全路径访问视图。这样当我们创建一个新的Session连接的时候，依旧可以访问全局的临时视图。

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdaQEKQeGCb3vnsJxy%2F-MFdemixQ9HTlGYMgsQt%2Fimage.png?alt=media\&token=e56d36d5-9c90-4a08-bd1b-a864955de06f)

通过全局临时视图，就可以跨越不同的Session进行访问。

## 3、DSL语法

&#x20;        DataFrame提供了一个特定领域语言(domain-specific language,DSL)去管理结构化数据，使用DSL语法就不必去创建临时视图了。我们也可以在Java、python、Scala和R中使用DSL。下面举例来说明：

### 1）创建一个DataFrame

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdWYAhB5fGc80pzSm1%2F-MFdWd9IeyZ7HwpiPGxO%2Fimage.png?alt=media\&token=ac0bc0f5-af3e-4f4d-8b49-6a357894fbbd)

依旧是通过读取json文件来创建一个DataFrame的对象。

### 2）查看DataFrame的Schema信息——**printSchema**&#x20;

这里我们使用DataFrame对象的 **printSchema** 方法进行查询

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdh1NocU-wNBCijZS8%2F-MFdhEDfuuO5SkF2jjza%2Fimage.png?alt=media\&token=5077d774-e5c4-4dc8-8011-e31265cfbc4d)

### 3）查看“name”列的数据以及"age+1"的数据——select

这里我们使用DataFrame对象的 **select** 方法进行查询

<div align="center"><img src="https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdi8SQZG-Go0VJ747Z%2F-MFdiDNMStgHYW5StItW%2Fimage.png?alt=media&amp;token=400d72d1-679f-4d12-a3f5-a8b33a5629bb" alt=""></div>

**如果涉及到运算操作的时候，每列都必须使用$或者是引号的形式表达(单引号+字段名)，例如年龄+1的操作。同时类似于数据库，我们可以将age+1操作后的结果从新起一个别名，则在其中加入as + "别名"。**

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdjejE2-6x-p7ldRSy%2F-MFdjnvrKWqWFaQCtkGe%2Fimage.png?alt=media\&token=5886765e-f40d-431c-9067-23dc9be6fa90)

### 4）查看年龄大于20的数据——filter

这里我们使用DataFrame对象的 **filter** 方法进行查询。

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdkXP_cxkzbwJ0n4EP%2F-MFdlnIZDF5SAUTbRbRV%2Fimage.png?alt=media\&token=eed64968-81af-4846-8ea0-6d6dfb21f5fa)

### 5）查看年龄分组的数据——groupBy

这里我们使用DataFrame对象的 **groupBy** 方法进行查询。

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdmGKG1UzR7RyuNpfB%2F-MFdmMPG3hoA6W5cUqQw%2Fimage.png?alt=media\&token=a9ea1e05-35bf-4c8e-9a2d-08a4d2ae7a46)

## 4、DataFrame与RDD的相互转换

在开始我们说过DataFrame的创建方式有三种，其中一种就是从现有的RDD进行转换，那么就看看如何将两者进行转换(注意：在IDEA开发时，将DF和RDD进行互操作时，必须要引入  import spark.implicits.\_     在spark-shell中则无需引入)

### 1）RDD转换成DataFrame

这里我们可以通过RDD对象中的 **toDF** 方法将RDD转换成DataFrame类型。

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdt63HAsyP7D-BoVBt%2F-MFdv4CJflCXRzuc0vdn%2Fimage.png?alt=media\&token=c22d8944-f680-4d22-bbed-e7c95c4aa72f)

上例先创建一个rdd对象，通过toDF方法，传递进去列名("id","name","age")，再将RDD转化成DataFrame对象。但是在实际工作中会使用 **样例类** 的方式进行转换。**这样做的好处是，在我们没有给toDF传入列名时，会使用样例类中的属性名来代替，如果此时依旧通过toDF传递列名(此时必须传入全部分列名)，则会覆盖样例类中的属性名。**

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdwV8bs4sVHzAYx171%2F-MFdwY37c7nUdm48MyVs%2Fimage.png?alt=media\&token=048b650f-edf7-4df4-8014-2a613f5dc94d)

### 2）DataFrame转换成RDD

由于DataFrame本身是RDD封装而来，所以可以直接通过DataFrame获取内部的RDD。

![](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdyQp8SVx5ecNeslYn%2F-MFdz3ARFkEzaWAn36ry%2Fimage.png?alt=media\&token=9cc1760e-bb17-46fd-b3e4-7e3db84d08e3)

此时通过rdd方法返回的是一个Row对象,对于Row对象其实就是一个Array数组，可以通过Array(index)进行访问。

总的来说DF与RDD的转换其实就是如下图所示：

![DF与RDD转换图](https://3317023386-files.gitbook.io/~/files/v0/b/gitbook-legacy-files/o/assets%2F-MEz5KNfZ0XQVN4cib9s%2F-MFdxVxdNgxTtWicKPfT%2F-MFdxdyMjnhW8WiqKf6H%2Fimage.png?alt=media\&token=2aff7338-6371-48a7-9808-a9781a26640a)

## 5、DataFrame 函数
