Skip to content

项目5:分析水稻品种审定数据——Spark SQL结构化数据文件处理 ​

先修基础:项目1-4(Spark概述、Scala基础、Spark Shell编程、Spark IDE编程)、Hive基础


目录 ​


第一部分:项目背景与 Spark SQL 概述 ​

1.1 项目背景 ​

为什么要学 Spark SQL? ​

前面我们学了 RDD 编程,功能很强,但有几个问题:

  • 写代码比较繁琐,简单的查询也要写一大串 map、filter、reduceByKey
  • 学过 SQL 的人上手慢,需要重新学 RDD API
  • 处理结构化数据(表格数据)不方便,没有"列"的概念

Spark SQL 就是来解决这些问题的:

  • 可以直接写 SQL 语句 查询数据,会 SQL 就能用
  • 提供了 DataFrame 数据模型,像操作表格一样操作数据
  • 性能比 RDD 更好(有 Catalyst 优化器、Tungsten 优化)

项目场景 ​

水稻是我国最重要的粮食作物之一,水稻品种的审定与推广关系到国家粮食安全。 现有一份水稻品种审定数据 ricedata.csv,包含品种名称、亲本来源、类型等 7 个字段。

本项目要做的事:

  1. 用 Spark SQL 读取 CSV 数据
  2. 探索和预处理数据(去重、去空值、去异常值)
  3. 统计分析(各省份审定数量、水稻类型分布、农业部审定类型情况)

1.2 RDD vs DataFrame vs SQL 对比 ​

对比项RDDDataFrameSQL
数据模型分布式对象集合分布式表格(有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 目录,所有节点都要复制。

bash
# 主节点复制
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 驱动。

bash
# 主节点复制
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 ​

bash
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 级别,看得清楚。

bash
# 复制模板文件
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步:分发配置文件到子节点 ​

bash
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步:启动集群和服务 ​

bash
# 启动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 命令行 ​

bash
cd $SPARK_HOME/bin
./spark-sql

进入 spark-sql 后就可以直接写 HiveQL 语句了:

sql
-- 查看数据库
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 就能用。

自动创建的变量 ​

变量类型说明
scSparkContextRDD 编程入口(之前学的)
sparkSparkSessionSpark SQL 入口(新的)

💡 注意:Spark 2.0 之前是 SQLContext 和 HiveContext,2.0 之后统一成 SparkSession 了。

在 spark-shell 中执行 SQL ​

scala
// 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。

基本创建方式 ​

scala
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 支持的创建方式 ​

scala
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 对比 ​

对比项SparkContextSparkSession
版本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 种方式:

  1. 从结构化数据文件创建(Parquet、JSON、CSV)
  2. 从外部数据库创建(MySQL、Oracle等,通过JDBC)
  3. 从 RDD 创建(反射推断 Schema)
  4. 从 RDD 创建(编程指定 Schema)
  5. 从 Hive 表创建

4.1 方式一:从结构化数据文件创建 ​

① Parquet 文件(默认格式) ​

Parquet 是 Spark SQL 默认的文件格式,列式存储,压缩率高。

scala
// load()默认就是Parquet格式
val dfUsers = spark.read.load("/tipdm/data/SparkSQL/users.parquet")

② JSON 文件 ​

scala
// 方式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 文件 ​

scala
// 方式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。

scala
// 设置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。

步骤 ​

  1. 定义 case class(样例类)
  2. 读取文件创建 RDD
  3. RDD 映射为 case class 对象
  4. 调用 toDF() 转成 DataFrame

代码示例 ​

scala
// 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。

步骤 ​

  1. 创建 RDD[Row]
  2. 用 StructType 定义 Schema
  3. 用 createDataFrame 应用 Schema

代码示例 ​

scala
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:整个表的结构,包含多个 StructField
  • StructField:一个字段的定义,包含字段名、类型、是否可空
  • 常用类型:StringType、IntegerType、DoubleType、LongType、BooleanType

4.5 方式五:从 Hive 表创建 ​

scala
// 切换数据库
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 的列名和类型。

scala
movies.printSchema()

输出类似:

root
 |-- movieId: integer (nullable = false)
 |-- title: string (nullable = true)
 |-- Genres: string (nullable = true)

💡 理解:就像 Excel 里看"表头",知道每列叫什么、是什么类型。


5.2 show():查看数据 ​

最常用的查看数据方法。

用法 ​

scala
// 默认显示前20行,最多显示20个字符
movies.show()

// 显示所有字符(不截断)
movies.show(false)

// 显示前5行
movies.show(5)

// 显示前5行,且不截断
movies.show(5, false)

参数说明 ​

参数类型默认值说明
numRowsInt20显示多少行
truncateBooleantrue是否截断长字符串

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

代码示例 ​

scala
// 获取第一行
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
scala
// 获取所有数据
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 查询有两种方式:

  1. SQL 方式:注册临时表,写 SQL 查询
  2. DSL 方式:直接在 DataFrame 上调用方法(select、where 等)

💡 理解:

  • SQL 方式:就像在 MySQL 里写查询语句
  • DSL 方式:就像用链式调用的方法查询

6.1 方式一:SQL 查询 ​

步骤 ​

  1. 把 DataFrame 注册成临时视图(表)
  2. 用 spark.sql() 写 SQL 查询

代码示例 ​

scala
// 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() ​

两者功能一样,都是按条件筛选。

scala
// 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() ​

scala
// 查指定列
val userSelect = user.select("userId", "gender")

// 用Column方式,可以做计算
import org.apache.spark.sql.functions.col
user.select(col("userId"), col("age") + 1).show()

③ 特殊处理字段:selectExpr() ​

可以直接写 SQL 表达式,支持别名、函数等。

scala
// 简单查询
user.selectExpr("userId", "age + 1 as newAge").show()

// 使用UDF函数
user.selectExpr("userId", "replaced(gender) as sex", "age").show()

💡 selectExpr 很方便,里面可以直接写 SQL 风格的表达式。


④ 自定义函数(UDF) ​

scala
// 注册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() ​

scala
// 获取zip字段(返回Column对象)
val userCol = user.col("zip")
val userApply = user("zip")  // 等价于apply

// 配合select使用
user.select(userCol).show()

⑥ 限制行数:limit() ​

注意:limit() 不是 Action,是转换操作,需要配合 show() 等才执行。

scala
// 取前3行
val userLimit = user.limit(3)
userLimit.show()

💡 limit vs take:

  • limit 是转换算子(懒执行),返回 DataFrame
  • take 是行动算子(立即执行),返回 Array[Row]

⑦ 排序:orderBy() / sort() ​

两者用法一样,默认升序。

scala
// 升序(默认)
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 对象。

scala
// 按性别分组
val userGroupBy = user.groupBy("gender")

GroupedData 常用方法 ​

方法说明
count()每组的数量
max(col)每组的最大值
min(col)每组的最小值
mean(col) / avg(col)每组的平均值
sum(col)每组的和
agg(...)多种聚合组合

示例 ​

scala
// 按性别分组,统计每组人数
user.groupBy("gender").count().show()

// 按性别分组,统计平均年龄
user.groupBy("gender").avg("age").show()

⑨ 连接:join() ​

连接两个 DataFrame,类似 SQL 的 JOIN。

基本用法 ​

scala
// 笛卡尔积(不推荐,数据量大会爆炸)
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() ​

不分组也能做聚合,或者在分组后做多种聚合。

scala
// 不分组,整体聚合
user.agg(min("age"), mean("occupation"), count("zip")).show()

// 分组后多种聚合
user.groupBy("gender").agg(
  min("age"),
  max("age"),
  avg("age"),
  count("*")
).show()

💡 agg 的好处:可以一次做多种聚合,不用多次计算。


⑪ 删除列:drop() ​

scala
// 删除一列
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() ​

新增一列,或者修改已有列。

scala
// 新增一列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() ​

scala
// 删除有空值的行
dfPeople.na.drop().show()

// 用指定值填充空值
dfPeople.na.fill(20).show()

// 按列填充不同值
dfPeople.na.fill(Map("age" -> 20, "name" -> "unknown")).show()

💡 空值处理策略:

  • 空值少 → 直接删除(drop)
  • 空值多 → 填充默认值(fill)

⑭ 统计摘要:describe() ​

给出数值列的统计信息(计数、均值、标准差、最小值、最大值)。

scala
// 统计age和gender列
user.describe("age", "gender").show()

输出类似:

+-------+------------------+
|summary|               age|
+-------+------------------+
|  count|              6040|
|   mean| 26.34337748344371|
| stddev|11.838943270786138|
|    min|                 1|
|    max|                56|
+-------+------------------+

💡 describe 就像 Excel 的数据分析工具,一眼看出数据分布。


⑮ 采样:sample() ​

从数据中抽样,用于快速观察或测试。

scala
// 不重复抽样,抽样率0.2
user.sample(false, 0.2).show(5)

// 重复抽样,抽样率0.2,指定随机种子
user.sample(true, 0.2, seed = 123).show(5)

参数说明 ​

参数类型说明
withReplacementBoolean是否重复抽样(true=有放回,false=无放回)
fractionDouble抽样比例(0~1)
seedLong随机种子(可选,相同种子结果相同)

💡 使用场景:海量数据先抽一点看看,或者开发测试用。


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() ​

基本用法 ​

scala
// 准备要保存的数据
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(默认)
  • json
  • csv
  • jdbc(写入数据库)
  • text

7.2 保存为表:saveAsTable() ​

保存成持久化的表,存在 Hive 元数据中,程序重启后还在。

scala
// 保存为表
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。

实现代码 ​

scala
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:探索与预处理数据 ​

需求 ​

  1. 检查重复记录
  2. 检查空值
  3. 检查异常值(亲本来源含"?"、审定编号含"/")
  4. 数据清洗(去重、去空、去异常值)

实现代码 ​

scala
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())

代码详解 ​

空值统计的写法 ​

scala
riceData.columns.map(colName => 
  sum(when(col(colName).isNull, 1).otherwise(0)).alias(colName)
): _*

这段代码的意思是:

  • 遍历所有列名
  • 对每一列,统计 null 的数量(用 when/otherwise 实现)
  • : _* 是把数组展开成可变参数传给 select()

💡 理解:就像对每一列都做一次"数空值"的操作,最后拼成一行结果。

异常值处理 ​

  • !col("亲本来源").contains("?"):亲本来源不包含问号的才保留
  • ! 表示"非",就是取反

任务5.3:统计分析数据 ​

需求1:统计省级以上部门审定的水稻数量 ​

按省份分组,统计各省份审定的水稻数量,按数量降序排列。

scala
riceDataCleaned
  .groupBy("省份")
  .count()
  .orderBy(-col("count"))
  .show(10)

需求2:统计不同水稻类型的占比情况 ​

统计各类型的数量,看分布情况。

scala
println("水稻类型的总数量:" + riceDataCleaned.select("类型").distinct().count())

riceDataCleaned
  .groupBy("类型")
  .count()
  .orderBy(-col("count"))
  .show()

需求3:统计农业农村部审定的水稻类型情况 ​

筛选出农业农村部审定的,再按类型分组统计。

scala
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,没指定正确的编码 解决:

scala
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 风格)
scala
// 正确写法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 合并分区
scala
df.coalesce(1).write.save("path")  // 合并成1个文件

问题8:save 报错文件已存在 ​

现象:报错 path already exists 原因:默认保存模式是 error,文件已存在就报错 解决:

scala
df.write.mode("overwrite").save("path")  // 覆盖模式

9.4 排错通用思路 ​

  1. 先看 printSchema:确认列名和类型对不对
  2. 先 show() 看数据:确认数据读取正确
  3. 小数据先测:先拿少量数据验证逻辑
  4. 看列名是否正确:列名写错是最常见的错误
  5. 注意类型:字符串和数字要区分开

第十部分:实习 / 面试高频考点 ​

10.1 概念类(高频) ​

Q1:Spark SQL 是什么?有什么特点? ​

Spark SQL 是 Spark 处理结构化数据的组件,可以看作分布式 SQL 查询引擎。 特点:

  1. 统一处理:支持 SQL 和 DataFrame/Dataset API
  2. 多数据源:Parquet、JSON、CSV、Hive、JDBC 等
  3. 性能优化:Catalyst 优化器、Tungsten 执行引擎
  4. 多语言支持:Scala、Java、Python、R

Q2:RDD 和 DataFrame 有什么区别? ​

  1. 数据模型:RDD 是分布式对象集合,DataFrame 是分布式表格(有Schema)
  2. 结构化:RDD 没有列的概念,DataFrame 有列名和类型
  3. 性能:DataFrame 有 Catalyst 优化,性能更好
  4. API:RDD 是函数式编程,DataFrame 是 DSL + SQL
  5. 适用场景: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 有哪些方式? ​

五种方式:

  1. 从结构化文件创建(Parquet、JSON、CSV)
  2. 从外部数据库创建(JDBC 连接 MySQL/Oracle 等)
  3. 从 RDD 创建(反射推断 Schema,用 case class)
  4. 从 RDD 创建(编程指定 Schema,用 StructType)
  5. 从 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 快? ​

  1. Catalyst 优化器:会优化查询计划,比如谓词下推、列裁剪
  2. Tungsten 执行引擎:优化内存使用和代码生成
  3. 有 Schema 信息,可以做更多优化
  4. 列式存储(Parquet 等)更高效

Q12:DataFrame 也是分布式的吗?分区方式和 RDD 一样吗? ​

是的,DataFrame 底层也是 RDD,所以也是分布式分区存储的。 可以用 df.rdd 拿到底层的 RDD。 分区概念和 RDD 一样,也有宽窄依赖、Shuffle 等。

Q13:Spark SQL 支持自定义函数吗? ​

支持,有三种 UDF:

  1. UDF(用户定义函数):一进一出
  2. UDAF(用户定义聚合函数):多进一出(类似 sum、count)
  3. UDTF(用户定义表生成函数):一进多出 简单的用 spark.udf.register 注册就能用。

10.3 实操类(高频) ​

Q14:怎么读取 CSV 文件?需要注意什么? ​

用 spark.read.csv() 常用 option:

  • header:是否有表头
  • inferSchema:是否自动推断类型
  • sep:分隔符
  • encoding:编码(中文文件特别注意) 大数据量建议手动指定 schema,不要用 inferSchema

Q15:怎么用 SQL 查询 DataFrame? ​

  1. 先注册临时视图:df.createOrReplaceTempView("表名")
  2. 然后用 spark.sql("SQL语句") 查询
  3. 返回的还是 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 文件,完成以下操作:

  1. 筛选出口味评分大于 7 分的数据
  2. 统计各类别餐饮点评数,并按降序排列
  3. 保存结果到 HDFS

实现代码 ​

scala
// 上传数据至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)按点评数总和降序排列
.rddDataFrame 转成 RDD
.repartition(1)合并成1个分区,输出1个文件
.saveAsTextFile()保存为文本文件

💡 注意:sum 后的列名会变成 sum(列名),排序时要用这个列名。


笔记版本:V1.0 对应教材:《Spark大数据技术与应用(第3版)》人民邮电出版社 对应项目:项目5 分析水稻品种审定数据——Spark SQL结构化数据文件处理 最后更新:2026年8月

基于 Vite 强力驱动 | 纯静态轻量托管