spark敲门砖之WordCount
Spark 入门实战:WordCount 程序全解析,从代码到执行
WordCount 是大数据的 Hello World——统计文本里每个单词出现几次。
Spark 版本的 WordCount 核心就四步:
- 读文件 → RDD[String]
- 拆单词 → flatMap
- 变键值对 → map(word => (word, 1))
- 聚合求和 → reduceByKey( + )
这四步写完,Spark 帮你分布式并行执行
环境准备:版本匹配与依赖配置
Spark 对 Scala 版本有严格依赖,错误的版本组合会导致兼容性问题(如类找不到、方法异常),需提前确认版本对应关系。
版本选择原则
- Spark 2.x 主要支持 Scala 2.11、2.12;
- Spark 3.x 主要支持 Scala 2.12、2.13(但部分早期 3.x 版本对 2.13 支持不完善);
- 推荐组合:Spark 3.1.1 + Scala 2.12.x(稳定性好,生态支持完善)。

Maven 依赖配置
在 pom.xml 中添加 Spark Core 和 Scala 依赖:
<!-- Scala 核心库 -->
<dependency>
<groupId>org.scala-lang</groupId>
<artifactId>scala-library</artifactId>
<version>2.12.13</version> <!-- 与 Spark 版本匹配 -->
</dependency>
<!-- Spark 核心依赖 -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.12</artifactId> <!-- _2.12 表示适配 Scala 2.12 -->
<version>3.1.1</version>
</dependency>
注意:
spark-coreartifactId 中的_2.12必须与 Scala 版本一致,否则会出现ClassNotFoundException。
代码:四步完成计数
import org.apache.spark.{SparkConf, SparkContext}
object WordCount {
def main(args: Array[String]): Unit = {
// 1. 配置
val conf = new SparkConf()
.setMaster("local[*]") // 本地模式,用所有 CPU 核心
.setAppName("WordCount")
// 2. 创建上下文(Spark 的"入口")
val sc = new SparkContext(conf)
// 3. 读文件 → RDD[String](每行一个元素)
val lines = sc.textFile("src/main/resources/wordcount.txt")
// 4. 拆单词 → flatMap(一行 → 多个单词)
val words = lines.flatMap(_.split(" "))
// 5. 变键值对 → map(每个单词 → (单词, 1))
val wordAndOne = words.map(word => (word, 1))
// 6. 聚合 → reduceByKey(相同 key 的 value 相加)
val wordCounts = wordAndOne.reduceByKey(_ + _)
// 7. 触发执行 + 打印结果
wordCounts.collect().foreach(println)
// 8. 关闭
sc.stop()
}
}
输入文件内容(wordcount.txt):
Hadoop Hive Spark
Hadoop HDFS
Hive Spark
Zookeeper Kafka Flume
输出:
(Hadoop, 2)
(Hive, 2)
(Spark, 2)
(HDFS, 1)
(Zookeeper, 1)
(Kafka, 1)
(Flume, 1)
关键概念:三步走
1. SparkConf + SparkContext = 启动入口
| 组件 | 做什么 |
|---|---|
SparkConf |
配置运行模式、应用名称、资源参数 |
SparkContext |
连接 Spark 集群,创建 RDD,调度任务 |
SparkContext 是所有 Spark 程序的”大门”——有了它才能操作 RDD。
2. RDD:分布式数据集合
textFile 读进来的 lines 是一个 RDD。它不是一个本地集合,是分布在多台机器上的数据分片。
RDD 的核心特性:
- 不可变:创建后不能改,要改就生成新 RDD
- 分区:数据分散在多台机器上
- 惰性:不遇到 Action 不会真的算
3. Transformation vs Action
| 类型 | 特点 | 例子 |
|---|---|---|
| Transformation | 惰性执行,只记录”要做什么”,返回新 RDD | flatMap、map、reduceByKey |
| Action | 触发真正执行,返回结果或写外部存储 | collect、count、saveAsTextFile |
代码里的执行时机:
lines.flatMap(...) ← 还没算,只是记下了
.map(...) ← 还没算,只是记下了
.reduceByKey(...) ← 还没算,只是记下了
.collect() ← 触发执行!所有前面的操作开始跑
这就是”惰性计算”——不等到要结果,绝不动手。
执行过程:4 个 Stage 在集群上怎么跑
textFile读文件 → 按 HDFS Block 或本地文件大小切分分区,每个分区一个 TaskflatMap+map→ 在每个分区上独立执行(窄依赖,不用 Shuffle)reduceByKey→ 触发 Shuffle,相同单词的数据聚到同一个分区做求和collect→ 把所有分区的结果拉到 Driver,合并成最终结果
窄依赖 vs 宽依赖:
| 依赖类型 | 特点 | 例子 |
|---|---|---|
| 窄依赖 | 父分区 → 子分区一对一,不用跨节点传输 | map、flatMap、filter |
| 宽依赖 | 父分区 → 子分区多对多,需要 Shuffle | reduceByKey、groupByKey |
窄依赖可以合并到一个 Stage,宽依赖必须切 Stage——Shuffle 是代价最大的操作。
Spark Shell 交互式运行(快速验证)
不用写完整程序,Spark 自带的 Shell 可以直接试:
spark-shell
进入后会自动创建好 sc(SparkContext),直接写:
val lines = sc.textFile("file:///path/to/wordcount.txt")
val counts = lines.flatMap(_.split(" ")).map((_, 1)).reduceByKey(_ + _)
counts.collect().foreach(println)
适合快速测试、验证逻辑、学习 API。