一、Spark概述
Apache Spark是一款快速通用的大数据处理引擎,支持批处理、流处理、机器学习和图计算等多种场景。Spark基于内存计算,相比Hadoop MapReduce有显著的性能优势。
二、Spark架构
2.1 Spark运行架构
graph TD
A[Driver] --> B[Cluster Manager]
B --> C[Worker节点1]
B --> D[Worker节点2]
B --> E[Worker节点3]
C --> C1[Executor]
D --> D1[Executor]
E --> E1[Executor]
C1 --> C2[Task]
C1 --> C3[Task]
D1 --> D2[Task]
D1 --> D3[Task]
2.2 Spark组件
| 组件 | 功能 | 说明 |
|---|---|---|
| Spark Core | 核心引擎 | RDD、任务调度 |
| Spark SQL | SQL处理 | DataFrame、Dataset |
| Spark Streaming | 流处理 | 微批处理 |
| MLlib | 机器学习 | 算法库 |
| GraphX | 图计算 | 图算法 |
三、RDD编程模型
3.1 RDD概念
RDD(Resilient Distributed Dataset)是Spark的核心数据抽象:
- 弹性:自动容错、数据重分区
- 分布式:数据分布在多个节点
- 数据集:不可变的分区集合
3.2 RDD操作
转换操作(Transformations)
// 创建RDD
val rdd = sc.parallelize(List(1, 2, 3, 4, 5))
// 转换操作
val mapped = rdd.map(x => x * 2)
val filtered = rdd.filter(x => x > 2)
val flatMapped = rdd.flatMap(x => List(x, x * 2))
// 聚合操作
val reduced = rdd.reduce((a, b) => a + b)
val grouped = rdd.groupBy(x => x % 2)
行动操作(Actions)
// 行动操作
rdd.count()
rdd.collect()
rdd.first()
rdd.take(3)
rdd.saveAsTextFile("output")
3.3 RDD持久化
// 持久化策略
rdd.cache() // 默认MEMORY_ONLY
rdd.persist(StorageLevel.MEMORY_ONLY)
rdd.persist(StorageLevel.MEMORY_AND_DISK)
rdd.persist(StorageLevel.DISK_ONLY)
四、DataFrame与Dataset
4.1 DataFrame
DataFrame是带有Schema的分布式数据集合:
// 创建DataFrame
val df = spark.read.json("data.json")
// 查看Schema
df.printSchema()
// 查询操作
df.select("name", "age").show()
df.filter(df("age") > 18).show()
df.groupBy("gender").count().show()
4.2 Dataset
Dataset是类型安全的DataFrame:
// 定义样例类
case class Person(name: String, age: Int)
// 创建Dataset
val ds = spark.read.json("data.json").as[Person]
// 类型安全操作
ds.filter(_.age > 18).show()
4.3 三种API对比
| API | 类型安全 | 性能 | 易用性 |
|---|---|---|---|
| RDD | 是 | 低 | 低 |
| DataFrame | 否 | 高 | 高 |
| Dataset | 是 | 高 | 高 |
五、Spark SQL
5.1 SQL查询
// 注册临时表
df.createOrReplaceTempView("people")
// 执行SQL查询
val result = spark.sql("SELECT name, age FROM people WHERE age > 18")
result.show()
5.2 数据源
// 读取CSV
val df = spark.read
.option("header", "true")
.csv("data.csv")
// 读取Parquet
val df = spark.read.parquet("data.parquet")
// 读取JDBC
val df = spark.read
.format("jdbc")
.option("url", "jdbc:mysql://localhost:3306/test")
.option("dbtable", "people")
.load()
六、Spark性能优化
6.1 数据分区
// 调整分区数
rdd.repartition(100)
rdd.coalesce(10)
// 设置默认分区数
spark.conf.set("spark.sql.shuffle.partitions", "200")
6.2 内存管理
// 内存配置
spark.driver.memory // Driver内存
spark.executor.memory // Executor内存
spark.executor.cores // Executor核心数
6.3 广播变量
// 使用广播变量
val broadcastVar = sc.broadcast(Array(1, 2, 3))
rdd.map(x => x * broadcastVar.value(0))
6.4 累加器
// 使用累加器
val accum = sc.longAccumulator("My Accumulator")
rdd.foreach(x => accum.add(x))
println(accum.value)
七、Spark与Hadoop对比
| 特性 | Spark | Hadoop MapReduce |
|---|---|---|
| 计算模型 | 内存计算 | 磁盘计算 |
| 性能 | 快10-100倍 | 较慢 |
| 编程模型 | 丰富 | 简单 |
| 数据格式 | 多种 | 主要文本 |
八、Spark批处理实践
8.1 ETL流程
flowchart TD
A[数据源] --> B[Spark读取]
B --> C[数据清洗]
C --> D[数据转换]
D --> E[数据聚合]
E --> F[数据写入]
F --> G[数据仓库]
8.2 数据清洗示例
val cleanedDF = rawDF
.filter(col("age").isNotNull)
.filter(col("email").rlike("^[A-Za-z0-9+_.-]+@[A-Za-z0-9.-]+$"))
.withColumn("age", col("age").cast("int"))
.dropDuplicates("id")
8.3 数据聚合示例
val resultDF = cleanedDF
.groupBy("country", "gender")
.agg(
count("id").alias("total"),
avg("age").alias("avg_age"),
sum("income").alias("total_income")
)
.orderBy(desc("total"))
九、Spark部署模式
| 部署模式 | 特点 | 适用场景 |
|---|---|---|
| Local模式 | 单机运行 | 开发测试 |
| Standalone模式 | Spark自带集群 | 中小型集群 |
| YARN模式 | 基于Hadoop YARN | 大型集群 |
| Kubernetes模式 | 基于K8s | 云原生环境 |
十、总结
Spark是大数据批处理的主流框架,掌握Spark的核心概念和编程模型,能够高效处理海量数据。通过合理的性能优化和部署配置,能够构建稳定可靠的大数据分析系统。