sparkSQL简介

Spark SQL:结构化数据在 Spark 上跑得更快、写得更爽

Spark 早期只有 RDD,什么数据都能处理——文本、JSON、二进制、自定义对象。

但处理结构化数据(比如 CSV、数据库表、JSON)的时候,RDD 用起来很麻烦:

// RDD 方式:手动解析 CSV
val rdd = sc.textFile("people.csv")
  .map(line => {
    val parts = line.split(",")
    (parts(0), parts(1).toInt)   // 手动解析、手动转类型
  })
  .filter(_._2 > 18)

每一行都要手动 split、手动转类型。写一次还行,写十次就烦了。

Spark SQL 就是为了解决这个问题:给结构化数据加上 schema(字段名 + 类型),然后用 SQL 或链式 API 操作,自动优化执行。

对比项 RDD Spark SQL(DataFrame)
有没有 schema 没有,数据就是对象 有,字段名+类型
怎么查 map、filter、reduce SQL 或 select、where、groupBy
执行优化 手写优化 Catalyst 优化器自动优化
代码量 多 少

Spark SQL 的核心组件

组件 做什么
Catalyst 优化器 自动优化 SQL 执行计划(谓词下推、列裁剪、常量折叠)
Tungsten 执行引擎 内存管理优化 + 代码生成,提升执行效率
DataFrame 有 schema 的分布式数据集(类比 SQL 表)
DataSet 强类型版本的 DataFrame(Scala/Java 专属)

Catalyst 做了你手动优化的事:

  • 只读需要的列(列裁剪)
  • 把过滤条件下推到数据源(谓词下推)
  • 合并可以合并的操作
-- SQL 写的是这样
SELECT name FROM people WHERE age > 18

-- Catalyst 优化后做的事:
-- 1. 只读 name 和 age 两列(不读其他列)
-- 2. 在数据源就过滤 age > 18(不把全量数据读进来再过滤)
-- 3. 只返回 name 列

DataFrame 和 DataSet 的关系

DataFrame DataSet
类型 无类型(DataFrame = Dataset[Row]) 强类型(Dataset[Person])
API SQL / 方法链 类型安全的方法链
编译时检查 没有 有
适用语言 Scala / Python / Java / R Scala / Java 为主
推荐度 日常分析首选 需要类型安全时用

Python 里只有 DataFrame(没有 DataSet),Scala 里两者都有。

Spark SQL 的使用流程

1. 创建 SparkSession(入口)

val spark = SparkSession.builder()
  .appName("MyApp")
  .master("local[*]")
  .getOrCreate()

import spark.implicits._   // 导入隐式转换

2. 读数据

// 读 JSON
val df = spark.read.json("people.json")

// 读 CSV
val df = spark.read
  .option("header", "true")
  .option("inferSchema", "true")
  .csv("people.csv")

// 读 Parquet
val df = spark.read.parquet("people.parquet")

3. 处理数据

// 方式一:方法链
df.select("name", "age")
  .where("age > 18")
  .groupBy("name")
  .count()
  .show()

// 方式二:SQL
df.createOrReplaceTempView("people")
spark.sql("SELECT name, COUNT(*) FROM people WHERE age > 18 GROUP BY name").show()

4. 写数据

df.write.parquet("output.parquet")
df.write.json("output.json")
df.write.jdbc(jdbcUrl, "table", props)

性能:DataFrame 比 RDD 快多少?

同样的数据、同样的操作,DataFrame 通常比 RDD 快 2-5 倍(复杂查询差距更大)。

原因:

  • Catalyst 优化器自动做谓词下推、列裁剪
  • Tungsten 引擎用二进制内存管理,减少 GC
  • 代码生成把表达式编译成 Java 字节码
// RDD 方式
rdd.map(_.age).filter(_ > 18).sum()   // 自己算

// DataFrame 方式
df.select("age").where("age > 18").rdd.map(_.getInt(0)).sum()
// 或者直接用 DataFrame 聚合
df.filter("age > 18").agg(sum("age"))
// DataFrame 方式有 Catalyst 优化

Spark SQL vs Hive

Spark SQL Hive
执行引擎 Spark(内存计算) MapReduce(磁盘计算)
速度 快(2-10 倍) 慢
数据源 多源(JSON/Parquet/JDBC) 主要是 HDFS
是否依赖 Hive 否(可选支持) 是
适用场景 交互式查询、ETL 离线数据仓库

Spark SQL 可以读 Hive 表:

val spark = SparkSession.builder()
  .enableHiveSupport()
  .getOrCreate()

spark.sql("SELECT * FROM hive_table").show()