项目5:分析水稻品种审定数据——Spark SQL结构化数据文件处理
先修基础:项目1-4(Spark概述、Scala基础、Spark Shell编程、Spark IDE编程)、Hive基础
目录
- 第一部分:项目背景与 Spark SQL 概述
- 第二部分:Spark SQL 基本概念与配置
- 第三部分:Spark SQL 与 Shell 交互
- 第四部分:DataFrame 创建方式(5种)
- 第五部分:DataFrame 查看数据操作
- 第六部分:DataFrame 查询操作(重点)
- 第七部分:DataFrame 输出操作
- 第八部分:项目实战——水稻品种审定数据分析
- 第九部分:常见问题与排错指南
- 第十部分:实习 / 面试高频考点
- 附录:习题解析
第一部分:项目背景与 Spark SQL 概述
1.1 项目背景
为什么要学 Spark SQL?
前面我们学了 RDD 编程,功能很强,但有几个问题:
- 写代码比较繁琐,简单的查询也要写一大串 map、filter、reduceByKey
- 学过 SQL 的人上手慢,需要重新学 RDD API
- 处理结构化数据(表格数据)不方便,没有"列"的概念
Spark SQL 就是来解决这些问题的:
- 可以直接写 SQL 语句 查询数据,会 SQL 就能用
- 提供了 DataFrame 数据模型,像操作表格一样操作数据
- 性能比 RDD 更好(有 Catalyst 优化器、Tungsten 优化)
项目场景
水稻是我国最重要的粮食作物之一,水稻品种的审定与推广关系到国家粮食安全。 现有一份水稻品种审定数据 ricedata.csv,包含品种名称、亲本来源、类型等 7 个字段。
本项目要做的事:
- 用 Spark SQL 读取 CSV 数据
- 探索和预处理数据(去重、去空值、去异常值)
- 统计分析(各省份审定数量、水稻类型分布、农业部审定类型情况)
1.2 RDD vs DataFrame vs SQL 对比
| 对比项 | RDD | DataFrame | SQL |
|---|---|---|---|
| 数据模型 | 分布式对象集合 | 分布式表格(有Schema) | 表格 |
| 是否有列名 | ❌ 没有 | ✅ 有列名和类型 | ✅ 有列名和类型 |
| 编程方式 | 函数式编程(map/filter等) | DSL API(select/where等) | SQL语句 |
| 学习门槛 | 高(要学Scala+函数式) | 中(类似表格操作) | 低(会SQL就行) |
| 性能优化 | 手动优化 | 自动优化(Catalyst) | 自动优化(Catalyst) |
| 适用场景 | 非结构化数据、复杂逻辑 | 结构化数据、数据分析 | 结构化数据、即席查询 |
💡 理解:
- RDD 就像"一堆散放的文件",你得自己整理
- DataFrame 就像"Excel 表格",有行列,可以按列操作
- SQL 就像"直接问问题",说你要什么就行
第二部分:Spark SQL 基本概念与配置
2.1 Spark SQL 是什么?
Spark SQL 是 Spark 四大组件之一,用于处理结构化数据。 可以把它理解为一个分布式的 SQL 查询引擎。
核心特点
- 统一处理:既能处理 RDD,也能处理结构化数据
- 多数据源支持:Parquet、JSON、CSV、Hive、JDBC 等
- 性能优化:Catalyst 优化器、Tungsten 执行引擎
- 多语言 API:Scala、Java、Python、R
发展历史
- Shark → Spark SQL(摆脱了对 Hive 的依赖)
- Spark 2.0 开始,入口统一为 SparkSession(之前是 SQLContext + HiveContext)
2.2 Spark SQL 配置(连接 Hive)
Spark SQL 可以独立运行,也可以连接 Hive 使用 Hive 的元数据和表。
配置步骤(7步)
第1步:复制 hive-site.xml
把 Hive 的配置文件复制到 Spark 的 conf 目录,所有节点都要复制。
# 主节点复制
cp /usr/local/hive-3.1.3/conf/hive-site.xml \
/usr/local/spark-3.5.1-bin-hadoop3-scala2.13/conf/
# 分发到子节点
scp /usr/local/hive-3.1.3/conf/hive-site.xml \
slave1:/usr/local/spark-3.5.1-bin-hadoop3-scala2.13/conf/
scp /usr/local/hive-3.1.3/conf/hive-site.xml \
slave2:/usr/local/spark-3.5.1-bin-hadoop3-scala2.13/conf/第2步:复制 MySQL 驱动包
Hive 元数据存在 MySQL 里,Spark 连接 Hive 需要 MySQL 驱动。
# 主节点复制
cp /usr/local/hive-3.1.3/lib/mysql-connector-java-8.0.30.jar \
/usr/local/spark-3.5.1-bin-hadoop3-scala2.13/jars/
# 分发到子节点
scp /usr/local/hive-3.1.3/lib/mysql-connector-java-8.0.30.jar \
slave1:/usr/local/spark-3.5.1-bin-hadoop3-scala2.13/jars/
scp /usr/local/hive-3.1.3/lib/mysql-connector-java-8.0.30.jar \
slave2:/usr/local/spark-3.5.1-bin-hadoop3-scala2.13/jars/第3步:配置 spark-env.sh
vim /usr/local/spark-3.5.1-bin-hadoop3-scala2.13/conf/spark-env.sh
# 添加以下内容
export SPARK_CLASSPATH=/usr/local/spark-3.5.1-bin-hadoop3-scala2.13/jars/mysql-connector-java-8.0.30.jar第4步:修改日志级别
默认日志太多,改成 warn 级别,看得清楚。
# 复制模板文件
cp $SPARK_HOME/conf/log4j2.properties.template $SPARK_HOME/conf/log4j2.properties
# 修改日志级别
vim $SPARK_HOME/conf/log4j2.properties
# 找到 rootLogger.level,改成 warn
rootLogger.level = warn第5步:分发配置文件到子节点
scp $SPARK_HOME/conf/spark-env.sh slave1:$SPARK_HOME/conf/
scp $SPARK_HOME/conf/spark-env.sh slave2:$SPARK_HOME/conf/
scp $SPARK_HOME/conf/log4j2.properties slave1:$SPARK_HOME/conf/
scp $SPARK_HOME/conf/log4j2.properties slave2:$SPARK_HOME/conf/第6步:启动集群和服务
# 启动Hadoop
/usr/local/hadoop-3.3.6/sbin/start-all.sh
# 启动Spark
/usr/local/spark-3.5.1-bin-hadoop3-scala2.13/sbin/start-all.sh
# 启动MySQL
systemctl start mysqld
# 启动Hive元数据服务
hive --service metastore &第7步:启动 spark-sql 命令行
cd $SPARK_HOME/bin
./spark-sql进入 spark-sql 后就可以直接写 HiveQL 语句了:
-- 查看数据库
show databases;
-- 创建表
create table students(
id int,
name string,
score double,
classes string
)
row format delimited fields terminated by '\t';第三部分:Spark SQL 与 Shell 交互
3.1 spark-shell 中的 Spark SQL
Spark SQL 已经集成在 spark-shell 里了,启动 spark-shell 就能用。
自动创建的变量
| 变量 | 类型 | 说明 |
|---|---|---|
sc | SparkContext | RDD 编程入口(之前学的) |
spark | SparkSession | Spark SQL 入口(新的) |
💡 注意:Spark 2.0 之前是 SQLContext 和 HiveContext,2.0 之后统一成 SparkSession 了。
在 spark-shell 中执行 SQL
// 1. 先把DataFrame注册成临时视图
df.createOrReplaceTempView("people")
// 2. 用spark.sql()执行SQL语句
val result = spark.sql("select name, age from people where age > 20")
// 3. 查看结果
result.show()3.2 IDEA 中创建 SparkSession
在 IDEA 中写 Spark SQL 程序,需要手动创建 SparkSession。
基本创建方式
import org.apache.spark.sql.SparkSession
object Test {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder() // 构建器模式
.appName("Test") // 应用名称
.master("local[*]") // 本地模式([*]表示用所有CPU核心)
.getOrCreate() // 获取或创建
}
}带 Hive 支持的创建方式
import org.apache.spark.sql.SparkSession
object Test {
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder()
.appName("Test")
.master("local[*]")
.enableHiveSupport() // 启用Hive支持
.getOrCreate()
}
}⚠️ 注意:启用 Hive 支持需要确保
hive-site.xml在 classpath 中(放在 resources 目录下)。
3.3 SparkSession vs SparkContext 对比
| 对比项 | SparkContext | SparkSession |
|---|---|---|
| 版本 | 1.x 就有 | 2.0 开始统一入口 |
| 用途 | RDD 编程 | Spark SQL、DataFrame、Dataset |
| 变量名 | sc(shell中) | spark(shell中) |
| 功能范围 | RDD 相关 | SQL + DataFrame + 也能拿SparkContext |
| 获取方式 | new SparkContext(conf) | SparkSession.builder().getOrCreate() |
💡 关系:SparkSession 里面包含了 SparkContext,可以通过
spark.sparkContext获取。
第四部分:DataFrame 创建方式(5种)
DataFrame 是 Spark SQL 的核心数据模型,可以理解为分布式的表格,有列名和类型(Schema)。
创建 DataFrame 有 5 种方式:
- 从结构化数据文件创建(Parquet、JSON、CSV)
- 从外部数据库创建(MySQL、Oracle等,通过JDBC)
- 从 RDD 创建(反射推断 Schema)
- 从 RDD 创建(编程指定 Schema)
- 从 Hive 表创建
4.1 方式一:从结构化数据文件创建
① Parquet 文件(默认格式)
Parquet 是 Spark SQL 默认的文件格式,列式存储,压缩率高。
// load()默认就是Parquet格式
val dfUsers = spark.read.load("/tipdm/data/SparkSQL/users.parquet")② JSON 文件
// 方式1:用format指定
val dfPeople = spark.read.format("json").load("/tipdm/data/SparkSQL/people.json")
// 方式2:直接用json()方法(更简洁)
val dfPeople = spark.read.json("/tipdm/data/SparkSQL/people.json")③ CSV 文件
// 方式1:用format指定
val df = spark.read.format("csv")
.option("header", "true") // 第一行是表头
.option("sep", ";") // 分隔符
.load("/tipdm/data/people.csv")
// 方式2:直接用csv()方法(更常用)
val dfPeople2 = spark.read
.option("header", "true") // 有表头
.option("inferSchema", "true") // 自动推断类型
.option("encoding", "GBK") // 编码
.csv("/tipdm/data/people.csv")CSV 常用 option 参数
| 参数 | 说明 | 示例 |
|---|---|---|
header | 第一行是否是表头 | true / false |
sep | 分隔符 | , / ; / \t |
inferSchema | 是否自动推断列类型 | true / false |
encoding | 文件编码 | UTF-8 / GBK |
nullValue | 空值表示 | "" / "NULL" |
💡 注意:
inferSchema会额外扫描一遍数据来推断类型,大数据量时建议手动指定 schema,更高效。
4.2 方式二:从外部数据库创建(JDBC)
可以从 MySQL、Oracle 等关系型数据库读取数据创建 DataFrame。
// 设置MySQL的URL
val url = "jdbc:mysql://master/test"
// 连接MySQL,读取people表
val jdbcDF = spark.read.format("jdbc")
.options(Map(
"url" -> url,
"user" -> "root",
"password" -> "123456",
"dbtable" -> "people"
))
.load()常用参数
| 参数 | 说明 |
|---|---|
url | 数据库连接地址 |
user | 用户名 |
password | 密码 |
dbtable | 表名(也可以是子查询) |
driver | 驱动类名(可选,自动识别) |
⚠️ 注意:需要确保数据库驱动包在 classpath 中。
4.3 方式三:从 RDD 创建(反射推断 Schema)
利用 Scala 的样例类(case class) 和反射机制,自动推断 Schema。
步骤
- 定义 case class(样例类)
- 读取文件创建 RDD
- RDD 映射为 case class 对象
- 调用
toDF()转成 DataFrame
代码示例
// 1. 定义样例类(字段名就是列名,字段类型就是列类型)
case class Person(name: String, age: Int)
// 2. 读取文件创建RDD
val data = sc.textFile("/tipdm/data/SparkSQL/people.txt")
.map(_.split(","))
// 3. RDD转成case class对象的RDD
val peopleRDD = data.map(p => Person(p(0), p(1).trim.toInt))
// 4. 转成DataFrame
val peopleDF = peopleRDD.toDF()⚠️ 注意:
- 只有 case class 才能被隐式转换为 DataFrame
- IDEA 中使用需要先导入隐式转换:
import spark.implicits._- case class 要定义在 main 方法外面
4.4 方式四:从 RDD 创建(编程指定 Schema)
当无法提前定义 case class 时(比如列名和类型是动态的),可以用编程方式指定 Schema。
步骤
- 创建 RDD[Row]
- 用 StructType 定义 Schema
- 用 createDataFrame 应用 Schema
代码示例
import org.apache.spark.sql.types._
import org.apache.spark.sql.Row
// 1. 创建RDD
val people = sc.textFile("/tipdm/data/SparkSQL/people.txt")
// 2. 定义Schema
val schemaString = "name age"
val schema = StructType(
schemaString.split(" ").map(fieldName =>
StructField(fieldName, StringType, nullable = true)
)
)
// 3. 创建RowRDD
val rowRDD = people.map(_.split(","))
.map(p => Row(p(0), p(1).trim))
// 4. 应用Schema,创建DataFrame
val peopleDataFrame = spark.createDataFrame(rowRDD, schema)Schema 定义说明
StructType:整个表的结构,包含多个 StructFieldStructField:一个字段的定义,包含字段名、类型、是否可空- 常用类型:
StringType、IntegerType、DoubleType、LongType、BooleanType
4.5 方式五:从 Hive 表创建
// 切换数据库
spark.sql("use test")
// 查询Hive表,返回DataFrame
val people = spark.sql("select * from people")
// 也可以直接写表名
val people2 = spark.table("test.people")💡 前提:要先
enableHiveSupport(),并且配置好了 hive-site.xml。
4.6 五种创建方式总结
| 创建方式 | 适用场景 | 难度 |
|---|---|---|
| 结构化文件 | Parquet/JSON/CSV等文件 | ⭐ 简单 |
| 外部数据库 | MySQL/Oracle等关系型数据库 | ⭐⭐ |
| RDD反射 | 有固定结构,能定义case class | ⭐⭐ |
| RDD编程指定Schema | 结构动态,无法提前定义case class | ⭐⭐⭐ |
| Hive表 | 数据已经在Hive里 | ⭐ 简单 |
第五部分:DataFrame 查看数据操作
DataFrame 也是惰性求值的,只有触发 Action 操作才会真正计算。
5.1 printSchema:打印数据模式
查看 DataFrame 的列名和类型。
movies.printSchema()输出类似:
root
|-- movieId: integer (nullable = false)
|-- title: string (nullable = true)
|-- Genres: string (nullable = true)💡 理解:就像 Excel 里看"表头",知道每列叫什么、是什么类型。
5.2 show():查看数据
最常用的查看数据方法。
用法
// 默认显示前20行,最多显示20个字符
movies.show()
// 显示所有字符(不截断)
movies.show(false)
// 显示前5行
movies.show(5)
// 显示前5行,且不截断
movies.show(5, false)参数说明
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
numRows | Int | 20 | 显示多少行 |
truncate | Boolean | true | 是否截断长字符串 |
5.3 first / head / take / takeAsList
获取若干行记录,返回 Row 对象。
| 方法 | 返回类型 | 说明 |
|---|---|---|
first() | Row | 获取第一行 |
head() | Row | 同 first() |
head(n: Int) | Array[Row] | 获取前 n 行 |
take(n: Int) | Array[Row] | 获取前 n 行 |
takeAsList(n: Int) | List[Row] | 获取前 n 行,返回 List |
代码示例
// 获取第一行
movies.first()
// 获取前3行
movies.head(3)
movies.take(3)
// 获取前3行,返回List
movies.takeAsList(3)💡 注意:
first()和head()功能一样;take()和takeAsList()都会把数据拉到 Driver 端。
5.4 collect / collectAsList
获取所有数据。
| 方法 | 返回类型 | 说明 |
|---|---|---|
collect() | Array[Row] | 获取所有数据,返回数组 |
collectAsList() | List[Row] | 获取所有数据,返回 List |
// 获取所有数据
movies.collect()
movies.collectAsList()⚠️ 警告:数据量大时千万不要用 collect()!会把所有数据拉到 Driver 端,导致 OOM(内存溢出)。 大数据量时用
show()看前几行,或者保存到文件。
5.5 查看数据方法对比
| 方法 | 返回类型 | 是否Action | 是否拉到Driver | 适用场景 |
|---|---|---|---|---|
printSchema() | Unit | ❌ | - | 看结构 |
show() | Unit | ✅ | 只拉显示的行 | 日常查看数据 |
first()/head() | Row | ✅ | 1行 | 看第一条 |
take(n) | Array[Row] | ✅ | n行 | 看前n条 |
collect() | Array[Row] | ✅ | 全部行 | 小数据量全量获取 |
第六部分:DataFrame 查询操作(重点)
DataFrame 查询有两种方式:
- SQL 方式:注册临时表,写 SQL 查询
- DSL 方式:直接在 DataFrame 上调用方法(select、where 等)
💡 理解:
- SQL 方式:就像在 MySQL 里写查询语句
- DSL 方式:就像用链式调用的方法查询
6.1 方式一:SQL 查询
步骤
- 把 DataFrame 注册成临时视图(表)
- 用
spark.sql()写 SQL 查询
代码示例
// 1. 注册临时视图
peopleDataFrame.createOrReplaceTempView("peopleTempTab")
// 2. 写SQL查询
val personsRDD = spark.sql("select name, age from peopleTempTab where age > 20")
// 3. 查看结果
personsRDD.show()临时视图 vs 全局临时视图
| 类型 | 方法 | 范围 |
|---|---|---|
| 临时视图 | createOrReplaceTempView | 只在当前 SparkSession 有效 |
| 全局临时视图 | createGlobalTempView | 跨 SparkSession 有效,查询时要加 global_temp. 前缀 |
6.2 方式二:DSL 查询(DataFrame API)
直接在 DataFrame 上调用方法,链式调用。
下面是常用的查询方法:
① 条件查询:where() / filter()
两者功能一样,都是按条件筛选。
// where方式
val userWhere = user.where("gender = 'F' and age = 18")
// filter方式
val userFilter = user.filter("gender = 'F' and age = 18")
// 也可以用Column表达式(类型更安全)
import org.apache.spark.sql.functions.col
val userFilter2 = user.filter(col("gender") === "F" && col("age") === 18)💡 注意:用 Column 表达式时,等于用
===,不等于用=!=,不是==。
② 查询指定字段:select()
// 查指定列
val userSelect = user.select("userId", "gender")
// 用Column方式,可以做计算
import org.apache.spark.sql.functions.col
user.select(col("userId"), col("age") + 1).show()③ 特殊处理字段:selectExpr()
可以直接写 SQL 表达式,支持别名、函数等。
// 简单查询
user.selectExpr("userId", "age + 1 as newAge").show()
// 使用UDF函数
user.selectExpr("userId", "replaced(gender) as sex", "age").show()💡 selectExpr 很方便,里面可以直接写 SQL 风格的表达式。
④ 自定义函数(UDF)
// 注册UDF
spark.udf.register("replaced", (x: String) => {
x match {
case "M" => 0
case "F" => 1
}
})
// 使用UDF
user.selectExpr("userId", "replaced(gender) as sex").show()⑤ 获取单个字段:col() / apply()
// 获取zip字段(返回Column对象)
val userCol = user.col("zip")
val userApply = user("zip") // 等价于apply
// 配合select使用
user.select(userCol).show()⑥ 限制行数:limit()
注意:limit() 不是 Action,是转换操作,需要配合 show() 等才执行。
// 取前3行
val userLimit = user.limit(3)
userLimit.show()💡 limit vs take:
- limit 是转换算子(懒执行),返回 DataFrame
- take 是行动算子(立即执行),返回 Array[Row]
⑦ 排序:orderBy() / sort()
两者用法一样,默认升序。
// 升序(默认)
user.orderBy("userId").show()
// 降序(3种写法)
import org.apache.spark.sql.functions.desc
user.orderBy(desc("userId")).show()
user.orderBy($"userId".desc).show()
user.orderBy(-user("userId")).show()
// sort用法一样
user.sort("userId").show()
user.sort(desc("userId")).show()💡 降序三种写法:
desc("列名")、$"列名".desc、-df("列名"),选一种自己喜欢的就行。
⑧ 分组:groupBy()
按指定字段分组,返回 GroupedData 对象。
// 按性别分组
val userGroupBy = user.groupBy("gender")GroupedData 常用方法
| 方法 | 说明 |
|---|---|
count() | 每组的数量 |
max(col) | 每组的最大值 |
min(col) | 每组的最小值 |
mean(col) / avg(col) | 每组的平均值 |
sum(col) | 每组的和 |
agg(...) | 多种聚合组合 |
示例
// 按性别分组,统计每组人数
user.groupBy("gender").count().show()
// 按性别分组,统计平均年龄
user.groupBy("gender").avg("age").show()⑨ 连接:join()
连接两个 DataFrame,类似 SQL 的 JOIN。
基本用法
// 笛卡尔积(不推荐,数据量大会爆炸)
val dfjoin = user.join(rating)
// 按指定字段连接(内连接,默认)
val dfJoin = user.join(rating, "userId")
// 按多个字段连接
val dfJoin2 = dfJoin.join(user, Seq("userId", "gender"), "left_outer")连接类型(joinType)
| 类型 | 说明 | 对应SQL |
|---|---|---|
inner | 内连接(默认) | INNER JOIN |
outer / full / full_outer | 全外连接 | FULL OUTER JOIN |
left_outer | 左外连接 | LEFT JOIN |
right_outer | 右外连接 | RIGHT JOIN |
left_semi | 左半连接 | LEFT SEMI JOIN |
💡 和 RDD 的 join 对比:
- RDD 的 join 只能按键连接,且必须是键值对 RDD
- DataFrame 的 join 更灵活,可以指定连接字段和连接类型,更像 SQL
⑩ 聚合:agg()
不分组也能做聚合,或者在分组后做多种聚合。
// 不分组,整体聚合
user.agg(min("age"), mean("occupation"), count("zip")).show()
// 分组后多种聚合
user.groupBy("gender").agg(
min("age"),
max("age"),
avg("age"),
count("*")
).show()💡 agg 的好处:可以一次做多种聚合,不用多次计算。
⑪ 删除列:drop()
// 删除一列
user.drop("zip").show(5)
// 删除多列
user.drop("occupation", "zip").show(5)
// 也可以传Column对象
import org.apache.spark.sql.functions.col
user.drop(col("occupation"), col("zip")).show(5)⑫ 增加列:withColumn()
新增一列,或者修改已有列。
// 新增一列newAge,值为age+5
import org.apache.spark.sql.functions.col
user.withColumn("newAge", col("age") + 5).show(5)
// 修改已有列(覆盖)
user.withColumn("age", col("age") + 1).show(5)⑬ 空值处理:na()
// 删除有空值的行
dfPeople.na.drop().show()
// 用指定值填充空值
dfPeople.na.fill(20).show()
// 按列填充不同值
dfPeople.na.fill(Map("age" -> 20, "name" -> "unknown")).show()💡 空值处理策略:
- 空值少 → 直接删除(drop)
- 空值多 → 填充默认值(fill)
⑭ 统计摘要:describe()
给出数值列的统计信息(计数、均值、标准差、最小值、最大值)。
// 统计age和gender列
user.describe("age", "gender").show()输出类似:
+-------+------------------+
|summary| age|
+-------+------------------+
| count| 6040|
| mean| 26.34337748344371|
| stddev|11.838943270786138|
| min| 1|
| max| 56|
+-------+------------------+💡 describe 就像 Excel 的数据分析工具,一眼看出数据分布。
⑮ 采样:sample()
从数据中抽样,用于快速观察或测试。
// 不重复抽样,抽样率0.2
user.sample(false, 0.2).show(5)
// 重复抽样,抽样率0.2,指定随机种子
user.sample(true, 0.2, seed = 123).show(5)参数说明
| 参数 | 类型 | 说明 |
|---|---|---|
withReplacement | Boolean | 是否重复抽样(true=有放回,false=无放回) |
fraction | Double | 抽样比例(0~1) |
seed | Long | 随机种子(可选,相同种子结果相同) |
💡 使用场景:海量数据先抽一点看看,或者开发测试用。
6.3 DataFrame 常用操作速查表
| 操作 | 方法 | 类似SQL |
|---|---|---|
| 条件筛选 | where() / filter() | WHERE |
| 选择列 | select() / selectExpr() | SELECT |
| 限制行数 | limit() | LIMIT |
| 排序 | orderBy() / sort() | ORDER BY |
| 分组 | groupBy() | GROUP BY |
| 连接 | join() | JOIN |
| 聚合 | agg() / count() / max() 等 | 聚合函数 |
| 删除列 | drop() | - |
| 新增列 | withColumn() | SELECT ... AS |
| 空值处理 | na.drop() / na.fill() | - |
| 统计摘要 | describe() | - |
| 采样 | sample() | - |
第七部分:DataFrame 输出操作
7.1 保存为文件:save()
基本用法
// 准备要保存的数据
val copyOfUser = user.select("userId", "gender", "age")
// 保存为JSON文件
copyOfUser.write
.format("json")
.mode("overwrite")
.option("header", "true")
.save("/tipdm/data/SparkSQL/copyOfUser.json")mode(保存模式)
| mode | 说明 |
|---|---|
overwrite | 覆盖(已存在就替换) |
append | 追加(已存在就追加) |
ignore | 忽略(已存在就不保存) |
error / errorifexists | 报错(默认,已存在就报错) |
支持的输出格式
parquet(默认)jsoncsvjdbc(写入数据库)text
7.2 保存为表:saveAsTable()
保存成持久化的表,存在 Hive 元数据中,程序重启后还在。
// 保存为表
copyOfUser.write.saveAsTable("copyUser")
// 查询表
spark.sql("select * from copyUser").show(5)
// 或者用table方法
spark.table("copyUser").show(5)内部表 vs 外部表
- 内部表(默认):表数据由元数据服务管理,删除表数据也删
- 外部表:数据在外部路径,删除表只删元数据,数据还在
💡 saveAsTable 的好处:
- 持久化,下次直接查表名就行
- 可以被其他 Spark 程序共享
- 有 Hive 元数据管理
第八部分:项目实战——水稻品种审定数据分析
8.0 数据说明
数据文件
ricedata.csv:水稻品种审定数据
字段说明(共7个字段)
| 序号 | 字段名 | 说明 |
|---|---|---|
| 1 | 品种名称 | 水稻品种的名称 |
| 2 | 亲本来源 | 品种的亲本信息 |
| 3 | 类型 | 水稻类型(籼稻、粳稻等) |
| 4 | 审定编号 | 审定的编号 |
| 5 | 省份 | 审定省份/部门 |
| 6 | 审定年份 | 审定的年份 |
| 7 | 品种来源 | 品种来源说明 |
编码
GBK 编码(注意读取时要指定 encoding)
任务5.1:获取数据(读取 CSV 创建 DataFrame)
需求
读取 HDFS 上的 ricedata.csv 文件,创建 DataFrame。
实现代码
import org.apache.spark.sql.SparkSession
object rice {
def main(args: Array[String]): Unit = {
// 1. 创建SparkSession
val spark = SparkSession.builder()
.appName("rice")
.master("local[*]")
.enableHiveSupport()
.getOrCreate()
// 2. 设置日志级别为WARN(减少日志输出)
spark.sparkContext.setLogLevel("WARN")
// 3. 读取CSV数据
val riceData = spark.read
.option("header", "true") // 第一行是表头
.option("inferSchema", "true") // 自动推断类型
.option("encoding", "GBK") // 编码是GBK
.csv("hdfs://master:8020/tipdm/data/SparkSQL/ricedata.csv")
// 4. 查看前5行
riceData.show(5)
}
}代码说明
| 代码 | 说明 |
|---|---|
.option("header", "true") | 第一行是列名 |
.option("inferSchema", "true") | 自动推断每列的数据类型 |
.option("encoding", "GBK") | 文件是GBK编码,不指定会乱码 |
spark.sparkContext.setLogLevel("WARN") | 设置日志级别,减少干扰信息 |
任务5.2:探索与预处理数据
需求
- 检查重复记录
- 检查空值
- 检查异常值(亲本来源含"?"、审定编号含"/")
- 数据清洗(去重、去空、去异常值)
实现代码
import org.apache.spark.sql.functions._
// ========== 1. 探索重复记录 ==========
println("去重前的数据总行数:" + riceData.count())
println("去重后的数据总行数:" + riceData.distinct().count())
// ========== 2. 探索各字段空值数量 ==========
riceData.select(riceData.columns.map(
colName => sum(when(col(colName).isNull, 1).otherwise(0)).alias(colName)
): _*).show()
// ========== 3. 探索异常值 ==========
// 亲本来源含"?"的记录
riceData.where(col("亲本来源").contains("?")).show(5, false)
println("亲本来源字段中包含'?'的记录:" +
riceData.where(col("亲本来源").contains("?")).count())
// 审定编号含"/"的记录
riceData.where(col("审定编号").contains("/")).show(5, false)
println("审定编号字段中包含'/'的记录:" +
riceData.where(col("审定编号").contains("/")).count())
// ========== 4. 数据预处理 ==========
val riceDataCleaned = riceData
.distinct() // 去重
.na.drop() // 删除空值行
.where(!col("亲本来源").contains("?")) // 删除含?的异常行
.where(!col("审定编号").contains("/")) // 删除含/的异常行
println("数据预处理后的数据总行数:" + riceDataCleaned.count())代码详解
空值统计的写法
riceData.columns.map(colName =>
sum(when(col(colName).isNull, 1).otherwise(0)).alias(colName)
): _*这段代码的意思是:
- 遍历所有列名
- 对每一列,统计 null 的数量(用 when/otherwise 实现)
: _*是把数组展开成可变参数传给 select()
💡 理解:就像对每一列都做一次"数空值"的操作,最后拼成一行结果。
异常值处理
!col("亲本来源").contains("?"):亲本来源不包含问号的才保留!表示"非",就是取反
任务5.3:统计分析数据
需求1:统计省级以上部门审定的水稻数量
按省份分组,统计各省份审定的水稻数量,按数量降序排列。
riceDataCleaned
.groupBy("省份")
.count()
.orderBy(-col("count"))
.show(10)需求2:统计不同水稻类型的占比情况
统计各类型的数量,看分布情况。
println("水稻类型的总数量:" + riceDataCleaned.select("类型").distinct().count())
riceDataCleaned
.groupBy("类型")
.count()
.orderBy(-col("count"))
.show()需求3:统计农业农村部审定的水稻类型情况
筛选出农业农村部审定的,再按类型分组统计。
riceDataCleaned
.where(col("省份") === "农业农村部")
.groupBy("类型")
.count()
.orderBy(-col("count"))
.show()代码说明
| 代码 | 说明 |
|---|---|
groupBy("省份").count() | 按省份分组计数 |
orderBy(-col("count")) | 按count降序排列(负号表示降序) |
col("省份") === "农业农村部" | 省份等于"农业农村部"(注意是===不是==) |
.distinct().count() | 去重后计数(有多少种不同的类型) |
第九部分:常见问题与排错指南
9.1 环境配置类问题
问题1:spark-sql 启动报错找不到 Hive 元数据
现象:启动 spark-sql 报 Metastore 相关错误 原因:
- Hive 的 metastore 服务没启动
- hive-site.xml 没复制到 Spark 的 conf 目录
- MySQL 驱动包没放对位置 解决:
- 确认
hive --service metastore &已启动 - 确认 hive-site.xml 在 Spark 的 conf 目录下
- 确认 mysql-connector-java jar 包在 Spark 的 jars 目录下
问题2:读取 CSV 文件中文乱码
现象:中文列名或数据显示成乱码 原因:文件编码不是 UTF-8,没指定正确的编码 解决:
spark.read.option("encoding", "GBK").csv("path")问题3:inferSchema 推断的类型不对
现象:数字列被推断成 String 类型 原因:列中有非数字的值,或者 inferSchema 推断不准确 解决:
- 小数据量可以用 inferSchema,大数据量建议手动指定 schema
- 用 StructType 手动定义每列的类型
9.2 DataFrame 操作类问题
问题4:列名用中文报错或有问题
现象:中文列名查询时报错 原因:编码问题,或者列名有特殊字符 解决:
- 确保读取时指定了正确的编码
- 用反引号包裹列名(SQL 方式)
- 用 col("列名") 方式(DSL 方式)
问题5:=== 和 == 搞混了
现象:filter 条件不生效,或者报错 原因:
==是 Scala 的等于比较,返回 Boolean===是 Spark SQL 的 Column 等于比较,返回 Column 解决:- 用 Column 表达式时,用
=== - 用字符串表达式时,用
=(SQL 风格)
// 正确写法1:字符串表达式
df.filter("age = 18")
// 正确写法2:Column表达式
df.filter(col("age") === 18)问题6:collect() 导致 Driver OOM
现象:程序报 OutOfMemoryError 原因:数据量太大,collect() 把所有数据拉到 Driver 端 解决:
- 大数据量不要用 collect()
- 用 show() 看前几行
- 用 saveAsTextFile 或 write.save 保存到文件
9.3 输出类问题
问题7:save 保存后有很多小文件
现象:输出目录下有很多 part-xxxxx 文件 原因:DataFrame 有多少个分区,就会输出多少个文件 解决:
- 保存前用 coalesce 或 repartition 合并分区
df.coalesce(1).write.save("path") // 合并成1个文件问题8:save 报错文件已存在
现象:报错 path already exists 原因:默认保存模式是 error,文件已存在就报错 解决:
df.write.mode("overwrite").save("path") // 覆盖模式9.4 排错通用思路
- 先看 printSchema:确认列名和类型对不对
- 先 show() 看数据:确认数据读取正确
- 小数据先测:先拿少量数据验证逻辑
- 看列名是否正确:列名写错是最常见的错误
- 注意类型:字符串和数字要区分开
第十部分:实习 / 面试高频考点
10.1 概念类(高频)
Q1:Spark SQL 是什么?有什么特点?
Spark SQL 是 Spark 处理结构化数据的组件,可以看作分布式 SQL 查询引擎。 特点:
- 统一处理:支持 SQL 和 DataFrame/Dataset API
- 多数据源:Parquet、JSON、CSV、Hive、JDBC 等
- 性能优化:Catalyst 优化器、Tungsten 执行引擎
- 多语言支持:Scala、Java、Python、R
Q2:RDD 和 DataFrame 有什么区别?
- 数据模型:RDD 是分布式对象集合,DataFrame 是分布式表格(有Schema)
- 结构化:RDD 没有列的概念,DataFrame 有列名和类型
- 性能:DataFrame 有 Catalyst 优化,性能更好
- API:RDD 是函数式编程,DataFrame 是 DSL + SQL
- 适用场景:RDD 适合非结构化数据和复杂逻辑,DataFrame 适合结构化数据分析
Q3:SparkSession 和 SparkContext 有什么区别?
- SparkContext 是 Spark 的老入口,用于 RDD 编程
- SparkSession 是 Spark 2.0 后的新统一入口,用于 Spark SQL、DataFrame
- SparkSession 里面包含了 SparkContext,可以通过 spark.sparkContext 获取
- spark-shell 中自动创建的变量:sc(SparkContext)、spark(SparkSession)
Q4:创建 DataFrame 有哪些方式?
五种方式:
- 从结构化文件创建(Parquet、JSON、CSV)
- 从外部数据库创建(JDBC 连接 MySQL/Oracle 等)
- 从 RDD 创建(反射推断 Schema,用 case class)
- 从 RDD 创建(编程指定 Schema,用 StructType)
- 从 Hive 表创建
Q5:DataFrame 怎么注册成表?有几种视图?
用 createOrReplaceTempView 注册临时视图 两种视图:
- 临时视图(createOrReplaceTempView):只在当前 SparkSession 有效
- 全局临时视图(createGlobalTempView):跨 SparkSession 有效,查询时加 global_temp. 前缀
Q6:DataFrame 的 save 和 saveAsTable 有什么区别?
- save:保存成文件(Parquet/JSON/CSV 等),就是普通文件
- saveAsTable:保存成表,在 Hive 元数据中注册,持久化,下次可以直接查表名
- saveAsTable 保存的表,程序重启后还在,只要连同一个元数据服务就行
Q7:where 和 filter 有什么区别?
功能一样,都是条件筛选,只是名字不同。 可以传字符串表达式,也可以传 Column 表达式。 用哪个都行,看个人习惯。
Q8:orderBy 和 sort 有什么区别?
功能一样,都是排序,默认升序。 降序可以用 desc("列名")、$"列名".desc、-df("列名") 等写法。
Q9:coalesce 和 repartition 在 DataFrame 中也能用吗?
可以,DataFrame 也有 coalesce 和 repartition 方法,用法和 RDD 类似。 coalesce 主要用于减少分区,repartition 可以增也可以减。 常用于输出时合并小文件。
Q10:怎么处理空值?
用 na 方法:
- na.drop():删除有空值的行
- na.fill(value):用指定值填充空值
- na.fill(Map("列名" -> 值)):按列填充不同值
10.2 原理类(中频)
Q11:Spark SQL 为什么比 RDD 快?
- Catalyst 优化器:会优化查询计划,比如谓词下推、列裁剪
- Tungsten 执行引擎:优化内存使用和代码生成
- 有 Schema 信息,可以做更多优化
- 列式存储(Parquet 等)更高效
Q12:DataFrame 也是分布式的吗?分区方式和 RDD 一样吗?
是的,DataFrame 底层也是 RDD,所以也是分布式分区存储的。 可以用 df.rdd 拿到底层的 RDD。 分区概念和 RDD 一样,也有宽窄依赖、Shuffle 等。
Q13:Spark SQL 支持自定义函数吗?
支持,有三种 UDF:
- UDF(用户定义函数):一进一出
- UDAF(用户定义聚合函数):多进一出(类似 sum、count)
- UDTF(用户定义表生成函数):一进多出 简单的用 spark.udf.register 注册就能用。
10.3 实操类(高频)
Q14:怎么读取 CSV 文件?需要注意什么?
用 spark.read.csv() 常用 option:
- header:是否有表头
- inferSchema:是否自动推断类型
- sep:分隔符
- encoding:编码(中文文件特别注意) 大数据量建议手动指定 schema,不要用 inferSchema
Q15:怎么用 SQL 查询 DataFrame?
- 先注册临时视图:df.createOrReplaceTempView("表名")
- 然后用 spark.sql("SQL语句") 查询
- 返回的还是 DataFrame,可以继续操作
Q16:groupBy 后可以做哪些聚合?
常用的:count、max、min、avg/mean、sum 多种聚合组合用 agg() 也可以用 agg 配合各种聚合函数
Q17:join 有哪些类型?
inner(内连接,默认)、outer/full(全外)、left_outer(左外)、right_outer(右外)、left_semi(左半) 和 SQL 的 JOIN 类型对应。
Q18:怎么把 DataFrame 转成 RDD?
用 df.rdd,得到的是 RDD[Row] 反过来,RDD 转 DataFrame 有反射和编程指定 Schema 两种方式。
附录:习题解析
选择题解析
1、答案:A 解析:Spark SQL 是一个用于处理结构化数据的框架,可被视为一个分布式的 SQL 查询引擎,提供了一个抽象的可编程数据模型 DataFrame。
2、答案:D 解析:collect() 方法用于查询 DataFrame 中所有的数据,并返回一个 Array 对象,该方法不接受任何参数。
3、答案:A 解析:Spark SQL 可通过 JDBC 连接或 ODBC 连接的方式访问外部数据库(如 MySQL、Oracle),从外部数据库中读取数据创建 DataFrame。
4、答案:D 解析:Spark SQL 可通过 saveAsTable() 方法将 DataFrame 对象保存/持久化到 Hive 表中。
5、答案:C 解析:对 DataFrame 对象使用 show() 方法可以查看 DataFrame 数据,默认显示前 20 条记录。
6、答案:B 解析:describe() 方法用于统计 DataFrame 的基础信息(如计数、均值、标准差等),返回一个新的 DataFrame;collect() 方法用于查询 DataFrame 中所有的数据,并返回一个 Array 对象;collectAsList() 方法和 collect() 方法类似,但返回的是 List 对象。
7、答案:D 解析: ① Spark SQL 的核心设计目标之一是无缝集成 SQL 查询与 Spark 程序,因此 Spark SQL 与 Spark 一样提供了 Java、Scala、Python、R 等语言的 API; ② Spark SQL 可以兼容 Hive 以便在 Spark SQL 中访问 Hive 表、使用 UDF 和使用 Hive 查询语言; ③ Spark SQL 提供了统一的数据访问接口,支持多种数据源(如 Hive、JDBC、Parquet、JSON、CSV 等); ④ Spark SQL 支持标准的数据连接方式,如通过 JDBC 连接或 ODBC 连接的方式访问外部数据库(如 MySQL、Oracle),与 Hive 元数据集成,实现元数据共享。
8、答案:B 解析:Spark SQL 中的 DataFrame 对象没有 rename() 方法,修改列名可以用 withColumnRenamed 或 selectExpr 取别名。
9、答案:B 解析:orderBy() 方法是根据指定字段进行排序,默认为升序排序。
10、答案:B 解析:df.select("col_name") 和 df.select(col("col_name")) 返回的都是只包含 col_name 列的 DataFrame 对象;df.columns 返回的是包含 df 中所有列名的字符串数组;df.col("col_name") 返回的是 df 中的 col_name 列(Column 对象)。
操作题解析
题目
读取 restaurant.csv 文件,完成以下操作:
- 筛选出口味评分大于 7 分的数据
- 统计各类别餐饮点评数,并按降序排列
- 保存结果到 HDFS
实现代码
// 上传数据至HDFS
hdfs dfs -put /opt/data/restaurant.csv /tipdm/data
// 进入spark-shell
spark-shell
// 1. 读取数据
val data = spark.read
.option("header", "true")
.option("inferSchema", "true")
.csv("/tipdm/data/restaurant.csv")
// 查看数据结构和前几行
data.printSchema()
data.show(5)
// 2. 筛选出口味评分大于7分的数据
val data1 = data.filter(col("口味") > 7)
data1.show()
// 3. 统计各类别餐饮点评数,按降序排列
val data2 = data.groupBy("类别")
.sum("点评数")
.orderBy(col("sum(点评数)").desc)
data2.show()
// 4. 保存结果至HDFS
data1.rdd.repartition(1).saveAsTextFile("/tipdm/data/restaurant/result1")
data2.rdd.repartition(1).saveAsTextFile("/tipdm/data/restaurant/result2")代码说明
| 代码 | 说明 |
|---|---|
filter(col("口味") > 7) | 筛选口味大于7分的 |
groupBy("类别").sum("点评数") | 按类别分组,求和点评数 |
orderBy(col("sum(点评数)").desc) | 按点评数总和降序排列 |
.rdd | DataFrame 转成 RDD |
.repartition(1) | 合并成1个分区,输出1个文件 |
.saveAsTextFile() | 保存为文本文件 |
💡 注意:sum 后的列名会变成
sum(列名),排序时要用这个列名。
笔记版本:V1.0 对应教材:《Spark大数据技术与应用(第3版)》人民邮电出版社 对应项目:项目5 分析水稻品种审定数据——Spark SQL结构化数据文件处理 最后更新:2026年8月