引言
大数据领域,Scala 是 Spark 的「母语」——Spark 核心用 Scala 写成,Spark 的 API 在 Scala 下最完整、最原生。理解 Scala,才能读懂 Spark 源码、写出高效的数据处理管道,也是进入大数据工程师岗位的硬门槛。
本文从「为什么是大数据,为什么是 Scala」讲起,系统梳理 Spark 的核心抽象(RDD/DataFrame/Dataset 三代演进)、transformations 与 actions 的延迟执行模型、宽窄依赖与 DAG 的调度原理,再到用 DataFrame 实战数据处理、优化手段(分区、缓存、广播),最后拓展到 Flink 与流批一体的 Scala 工程范式。
前置:Scala 基础与集合操作(https://plumephp.com/scala-functional-programming/、https://plumephp.com/scala-collections/)。大数据分布式理论可参考 [[distributed-systems]] 专题。
目录
- 1. 大数据问题与 Spark 的答案
- 2. RDD:不可变分布式数据集
- 3. transformations 与 actions:延迟执行模型
- 4. DataFrame 与 Dataset:结构化数据主流
- 5. 宽窄依赖与 DAG:Spark 如何调度
- 6. Spark SQL 实战:从数据到洞察
- 7. 性能优化:分区、缓存与广播
- 8. Flink 与流批一体:Scala 在大数据的延伸
- 9. 总结:Scala 大数据开发的工程心智
- 延伸阅读
1. 大数据问题与 Spark 的答案
1.1 问题:单机内存装不下、算不动
| 问题 | 表现 |
|---|---|
| 数据量 | TB/PB 级,单机不行 |
| 计算量 | 复杂聚合,单机太慢 |
| 容错 | 机器随时会挂 |
1.2 Spark 的答案
- 分布式存储:数据分片放多机。
- 并行计算:任务分到各节点同时跑。
- 内存计算:中间结果驻留内存,比 Hadoop 快。
- 容错:血缘(lineage)重算丢失分区。
1.3 为什么用 Scala 写
- Spark 内核 Scala,API 原生。
- 函数式风格(map/filter)与数据处理高度契合。
- 类型安全 + JVM 生态。
2. RDD:不可变分布式数据集
2.1 RDD 是什么
RDD(Resilient Distributed Dataset)是不可变的分布式数据集合。它不真存数据,而是描述「如何从源头计算出来」(血缘)。
import org.apache.spark.{SparkConf, SparkContext}
val conf = new SparkConf().setAppName("demo").setMaster("local[*]")
val sc = new SparkContext(conf)
val rdd = sc.parallelize(1 to 1000) // 造一个 RDD
val sum = rdd.map(_ * 2).reduce(_ + _) // 变换 + 动作
println(sum) // 1001000
2.2 三代抽象演进
| 抽象 | 类型 | 场景 |
|---|---|---|
| RDD | 强类型、函数式 | 底层、自定义 |
| DataFrame | 弱类型(Row)、SQL | 结构化分析 |
| Dataset | 强类型 + DataFrame | 需要类型安全 |
3. transformations 与 actions:延迟执行模型
3.1 两类操作
| 类型 | 例子 | 何时执行 |
|---|---|---|
| transformation | map/filter/join/groupBy | 延迟,构建 DAG |
| action | collect/count/reduce/save | 触发真正计算 |
val rdd = sc.parallelize(1 to 100)
val filtered = rdd.filter(_ % 2 == 0) // transformation,不执行
val count = filtered.count() // action,此时才算
3.2 为什么延迟执行
- 合并优化:把多个变换合并成一个 task。
- 只算需要:action 只触发依赖链。
3.3 常用 transformations
rdd.map(_ * 2) // 逐个变换
rdd.flatMap(x => Seq(x, x)) // 展开
rdd.filter(_ > 10) // 过滤
rdd.distinct() // 去重
rdd.groupByKey() // 按键分组
rdd.reduceByKey(_ + _) // 按键聚合(比 groupByKey 高效)
rdd.join(other) // 连接
4. DataFrame 与 Dataset:结构化数据主流
4.1 创建 DataFrame
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder().appName("demo").master("local[*]").getOrCreate()
import spark.implicits._
val df = spark.read.option("header", true).csv("data/users.csv")
df.show(5)
df.printSchema()
4.2 DataFrame 操作
import org.apache.spark.sql.functions._
val result = df
.filter(col("age") > 18)
.groupBy("city")
.agg(count("*").as("cnt"), avg("age").as("avg_age"))
.orderBy(desc("cnt"))
result.show()
4.3 Dataset:带类型的 DataFrame
case class User(name: String, age: Int)
val ds: Dataset[User] = spark.read.json("data/users.json").as[User]
val adults = ds.filter(_.age > 18) // 类型安全,编译期检查
5. 宽窄依赖与 DAG:Spark 如何调度
5.1 依赖类型
| 依赖 | 含义 | 例子 |
|---|---|---|
| 窄依赖 | 一个父分区对应一个子分区 | map/filter |
| 宽依赖 | 父分区对多个子分区(需 shuffle) | groupByKey/join |
5.2 DAG 与 Stage
RDD 血缘构成 DAG → 按宽依赖切分 Stage → 每个 Stage 内任务并行执行
5.3 为什么宽依赖昂贵
宽依赖要 shuffle:跨节点搬运数据、落盘、网络传输。减少 shuffle 是 Spark 调优的核心。
6. Spark SQL 实战:从数据到洞察
6.1 注册临时表用 SQL
df.createOrReplaceTempView("users")
val sqlDf = spark.sql("""
SELECT city, COUNT(*) AS cnt, AVG(age) AS avg_age
FROM users
WHERE age > 18
GROUP BY city
ORDER BY cnt DESC
""")
sqlDf.show()
6.2 读写外部存储
// 读
val fromParquet = spark.read.parquet("s3://bucket/data/part-*")
// 写
df.write.mode("overwrite").partitionBy("city").parquet("out/users")
6.3 常见聚合场景
df.groupBy("city").count().show()
df.rolling(...) // 窗口(时间序列)
df.withColumn("ratio", col("a") / col("b")).show()
7. 性能优化:分区、缓存与广播
7.1 分区与并行度
// 调整分区数控制并行度
val repartitioned = df.repartition(200) // 增分区(可能 shuffle)
val coalesced = df.coalesce(20) // 减分区(尽量不 shuffle)
7.2 缓存:重复使用中间结果
val cached = df.cache() // 缓存在内存
// 或
df.persist(StorageLevel.MEMORY_AND_DISK)
7.3 广播:小表广播避免大 shuffle
import org.apache.spark.sql.functions.broadcast
val smallDf = spark.read.csv("data/dim.csv") // 小维表
val joined = bigDf.join(broadcast(smallDf), "key") // 广播小表,避免 shuffle
7.4 调优速查
| 问题 | 手段 |
|---|---|
| shuffle 大 | 广播小表、减少 join |
| 数据倾斜 | 加盐打散、salting |
| 重复计算 | cache/persist |
| 并行不够 | 增分区 |
8. Flink 与流批一体:Scala 在大数据的延伸
8.1 Flink 的 Scala API
import org.apache.flink.streaming.api.scala._
val env = StreamExecutionEnvironment.getExecutionEnvironment
val text = env.socketTextStream("host", 9999)
text
.flatMap(_.split(" "))
.map((_, 1))
.keyBy(_._1)
.sum(1)
.print()
env.execute("word count")
8.2 流批一体范式
| 框架 | 定位 | Scala 支持 |
|---|---|---|
| Spark | 批处理为主,微批流 | 原生 |
| Flink | 流处理原生,批流一体 | 原生 |
8.3 Scala 大数据工作流
数据接入(Kafka) → 流处理(Flink/Spark Streaming) → 批处理(Spark) → 数仓(SQL) → 分析/可视化
9. 总结:Scala 大数据开发的工程心智
9.1 三句话记住
- Spark 是「懒 + 分布式」:transformations 攒着,action 才跑。
- shuffle 是性能敌人:能广播就广播,能减就减。
- Scala 函数式恰好契合:map/filter/聚合即数据处理。
9.2 工程建议
| 建议 | 理由 |
|---|---|
| DataFrame 优先 | 引擎优化更好 |
| 减少宽依赖 | 降低 shuffle 成本 |
| 小表广播 | 避免大 join |
| 类型安全用 Dataset | 编译期查错 |
延伸阅读
- https://plumephp.com/scala-collections/ — 集合操作是大数据处理的单机版
- https://plumephp.com/scala-type-system/ — Dataset 类型安全的底层
- [[distributed-systems]] 专题 — 分布式系统理论
- Spark 官方文档 与 Flink 文档
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。