📖 数据密集型设计

Spark批处理与大数据分析

深入探讨Spark批处理框架与大数据分析实践

一、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的核心概念和编程模型,能够高效处理海量数据。通过合理的性能优化和部署配置,能够构建稳定可靠的大数据分析系统。