spark敲门砖之WordCount

Spark 入门实战:WordCount 程序全解析,从代码到执行

WordCount 是大数据的 Hello World——统计文本里每个单词出现几次。

Spark 版本的 WordCount 核心就四步:

  1. 读文件 → RDD[String]
  2. 拆单词 → flatMap
  3. 变键值对 → map(word => (word, 1))
  4. 聚合求和 → 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版本对应

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-core artifactId 中的 _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 flatMapmapreduceByKey
Action 触发真正执行,返回结果或写外部存储 collectcountsaveAsTextFile

代码里的执行时机:

lines.flatMap(...)     ← 还没算,只是记下了
  .map(...)            ← 还没算,只是记下了
  .reduceByKey(...)    ← 还没算,只是记下了
  .collect()           ← 触发执行!所有前面的操作开始跑

这就是”惰性计算”——不等到要结果,绝不动手。

执行过程:4 个 Stage 在集群上怎么跑

  1. textFile 读文件 → 按 HDFS Block 或本地文件大小切分分区,每个分区一个 Task
  2. flatMap + map → 在每个分区上独立执行(窄依赖,不用 Shuffle)
  3. reduceByKey → 触发 Shuffle,相同单词的数据聚到同一个分区做求和
  4. collect → 把所有分区的结果拉到 Driver,合并成最终结果

窄依赖 vs 宽依赖:

依赖类型 特点 例子
窄依赖 父分区 → 子分区一对一,不用跨节点传输 mapflatMapfilter
宽依赖 父分区 → 子分区多对多,需要 Shuffle reduceByKeygroupByKey

窄依赖可以合并到一个 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。