Spark 3.3.0 词频统计全流程实操指南(Scala 版,适配 Hadoop 集群)在这里插入图片描述

一、环境与配置说明

1.1 集群环境参数

项目配置详情说明
集群拓扑Master:hadoop102(192.168.10.102);Worker:hadoop103、hadoop1043 节点 Spark Standalone 集群
软件版本JDK 1.8.0_212;Hadoop 3.1.3;Spark 3.3.0;Maven 3.6.3已验证兼容性,支持分布式词频计算
核心环境变量export JAVA_HOME=/opt/module/jdk1.8.0_212``export HADOOP_HOME=/opt/module/hadoop-3.1.3``export SPARK_HOME=/opt/module/spark-3.3.0``export SPARK_MASTER_HOST=192.168.10.102``export MAVEN_HOME=/opt/module/maven-3.6.3``export PATH=$MAVEN_HOME/bin:$SPARK_HOME/bin:$HADOOP_HOME/bin:$PATH配置于所有节点/etc/profile

1.2 前置检查(必须执行)

步骤 1:检查 Spark 集群状态

在 Master 节点(hadoop102)启动并验证集群:

# 启动 Spark 集群(Master + Worker)
$SPARK_HOME/sbin/start-all.sh

# 验证 Master 节点进程(hadoop102 执行)
jps  # 需显示 "Master" 进程

# 验证 Worker 节点进程(hadoop103、hadoop104 分别执行)
jps  # 需显示 "Worker" 进程
步骤 2:验证 HDFS 服务可用性
# 查看 HDFS 根目录
hdfs dfs -ls /

# 测试 HDFS 读写权限(创建并删除临时文件)
hdfs dfs -touchz /test.txt && hdfs dfs -rm /test.txt

二、数据准备(HDFS 数据上传)

2.1 创建 HDFS 数据目录

# 创建输入目录(存放待分析文本)
hdfs dfs -mkdir -p /spark-wordcount/input

# 创建输出目录(仅演示,运行前需删除)
hdfs dfs -mkdir -p /spark-wordcount/output

2.2 准备本地测试文本

# 创建本地文本文件
vi /home/atguigu/wordcount-data.txt

粘贴测试内容:

Hello Spark Hello Hadoop
Spark is fast Hadoop is stable
Hello Spark HBase
Spark Spark Spark

2.3 上传文本到 HDFS

# 上传本地文件到 HDFS 输入目录
hdfs dfs -put /home/atguigu/wordcount-data.txt /spark-wordcount/input/

# 验证上传结果(需显示 wordcount-data.txt)
hdfs dfs -ls /spark-wordcount/input/

三、Scala 词频统计实现

3.1 编写 Scala 代码(目录结构严格遵循 Maven 规范)

步骤 1:创建代码目录(关键:符合 Maven 源码结构)
# 创建项目根目录
mkdir -p /home/atguigu/spark-code/scala

# 进入项目根目录
cd /home/atguigu/spark-code/scala

# 创建 Maven 规定的 Scala 源码目录(必须,否则 Maven 找不到代码)
mkdir -p src/main/scala
步骤 2:编写 Scala 代码
# 在 src/main/scala 目录下创建代码文件
vi src/main/scala/WordCount.scala
// 导入 Spark 核心依赖
import org.apache.spark.{SparkConf, SparkContext}

// 词频统计主类(单例对象,Scala 入口类)
object WordCount {
  def main(args: Array[String]): Unit = {
    // 1. 配置 Spark 应用
    val conf = new SparkConf()
      .setAppName("Scala-WordCount")  // 应用名称(Web UI 可见)
      .set("spark.executor.memory", "1g")  // 每个 Executor 内存
      .set("spark.cores.max", "4")  // 最大使用核心数

    // 2. 创建 Spark 上下文(核心入口)
    val sc = new SparkContext(conf)
    sc.setLogLevel("WARN")  // 减少日志输出

    // 3. 处理输入输出路径(优先使用命令行参数)
    val inputPath = if (args.length > 0) args(0) else "/spark-wordcount/input"
    val outputPath = if (args.length > 1) args(1) else "/spark-wordcount/output"

    // 4. 核心词频统计逻辑
    val wordCounts = sc.textFile(inputPath)  // 读取文件
      .flatMap(line => line.split("\\s+"))  // 按空白分词(扁平化)
      .map(word => (word, 1))  // 映射为(单词,1)
      .reduceByKey(_ + _)  // 按单词聚合计数
      .sortBy(_._2, ascending = false)  // 按词频降序排序

    // 5. 输出结果(控制台 + HDFS)
    println("=== 词频统计结果 ===")
    wordCounts.collect().foreach(println)  // 控制台打印(小数据)
    wordCounts.saveAsTextFile(outputPath)  // 保存到 HDFS

    // 6. 释放资源
    sc.stop()
  }
}
步骤 3:验证代码位置
# 确保代码在正确目录(src/main/scala)
ls -l src/main/scala/WordCount.scala
# 预期输出:-rw-r--r-- 1 atguigu atguigu ... src/main/scala/WordCount.scala

3.2 打包 Scala 代码(Maven 方式)

步骤 1:创建 pom.xml 配置文件
# 在项目根目录创建 pom.xml
vi /home/atguigu/spark-code/scala/pom.xml

粘贴配置(适配 Spark 3.3.0 + Scala 2.12):

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <!-- 项目基本信息 -->
    <groupId>com.atguigu.spark</groupId>
    <artifactId>spark-wordcount</artifactId>
    <version>1.0-SNAPSHOT</version>

    <!-- 版本管理(Scala 与 Spark 版本必须匹配) -->
    <properties>
        <scala.version>2.12.15</scala.version>
        <spark.version>3.3.0</spark.version>
        <maven.compiler.source>1.8</maven.compiler.source>
        <maven.compiler.target>1.8</maven.compiler.target>
    </properties>

    <!-- 依赖配置 -->
    <dependencies>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.12</artifactId>
            <version>${spark.version}</version>
            <scope>provided</scope>  <!-- 集群已存在,打包不包含 -->
        </dependency>
    </dependencies>

    <!-- 构建配置(编译 + 打包) -->
    <build>
        <sourceDirectory>src/main/scala</sourceDirectory>  <!-- Scala 源码目录 -->
        <plugins>
            <!-- Scala 编译插件 -->
            <plugin>
                <groupId>net.alchim31.maven</groupId>
                <artifactId>scala-maven-plugin</artifactId>
                <version>4.5.4</version>
                <executions>
                    <execution>
                        <goals>
                            <goal>compile</goal>
                            <goal>testCompile</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>

            <!-- 打包插件(生成可执行 JAR) -->
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-shade-plugin</artifactId>
                <version>3.3.0</version>
                <executions>
                    <execution>
                        <phase>package</phase>
                        <goals><goal>shade</goal></goals>
                        <configuration>
                            <filters>
                                <filter>
                                    <artifact>*:*</artifact>
                                    <excludes>
                                        <exclude>META-INF/*.SF</exclude>
                                        <exclude>META-INF/*.DSA</exclude>
                                        <exclude>META-INF/*.RSA</exclude>
                                    </excludes>
                                </filter>
                            </filters>
                            <!-- 指定主类(与代码中 object 名称一致) -->
                            <transformers>
                                <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
                                    <mainClass>WordCount</mainClass>
                                </transformer>
                            </transformers>
                        </configuration>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>
</project>
步骤 2:编译并打包(含关键验证)
# 进入项目根目录
cd /home/atguigu/spark-code/scala

# 清理旧文件并编译
mvn clean scala:compile

# 【关键验证 1】检查是否生成 .class 文件(必须成功)
ls -l target/classes/WordCount.class
# 预期输出:-rw-r--r-- 1 atguigu atguigu ... WordCount.class

# 打包生成 JAR
mvn package

# 【关键验证 2】检查 JAR 中是否包含主类
jar tf target/spark-wordcount-1.0-SNAPSHOT.jar | grep "WordCount.class"
# 预期输出:WordCount.class

四、提交任务到 Spark 集群

4.1 上传 JAR 到 HDFS

# 创建 HDFS 存放 JAR 的目录(若不存在)
hdfs dfs -mkdir -p /spark-jars/

# 上传本地 JAR 到 HDFS(-f 强制覆盖旧文件)
hdfs dfs -put -f target/spark-wordcount-1.0-SNAPSHOT.jar /spark-jars/

4.2 删除旧输出目录(避免冲突)

hdfs dfs -rm -r /spark-wordcount/output

4.3 提交任务

$SPARK_HOME/bin/spark-submit \
--master spark://192.168.10.102:7077 \
--class WordCount \
hdfs://192.168.10.102:8020/spark-jars/spark-wordcount-1.0-SNAPSHOT.jar \
/spark-wordcount/input \
/spark-wordcount/output

五、查看词频统计结果

5.1 查看 HDFS 结果文件

# 查看输出目录结构
hdfs dfs -ls /spark-wordcount/output

# 查看词频结果(part-00000 为结果文件)
hdfs dfs -cat /spark-wordcount/output/part-00000

预期输出

(Spark,6)
(Hello,3)
(Hadoop,2)
(is,2)
(fast,1)
(stable,1)
(HBase,1)

5.2 通过 Spark Web UI 监控(可选)

  • 任务运行时:访问 http://192.168.10.102:4040 查看实时进度;
  • 任务结束后:访问 http://192.168.10.102:8080 查看历史记录。

六、常见问题与解决方案

6.1 Maven 版本过低(报错:requires Maven version 3.3.9)

解决:升级 Maven 到 3.6.3(参考环境配置部分的 MAVEN_HOME 配置)。

6.2 主类找不到(ClassNotFoundException: WordCount)

解决

  1. 确保代码在 src/main/scala 目录;
  2. 验证 pom.xml<mainClass>WordCount</mainClass> 与代码类名一致;
  3. 重新打包并覆盖 HDFS 上的 JAR。

6.3 结果缺失(如缺少 Hello、Spark)

解决

  1. 验证 HDFS 输入文件内容:hdfs dfs -cat /spark-wordcount/input/wordcount-data.txt
  2. 检查提交命令的输入路径是否正确(需指向 /spark-wordcount/input)。

6.4 输出目录已存在(FileAlreadyExistsException)

解决:提交前删除旧输出目录:hdfs dfs -rm -r /spark-wordcount/output

七、流程总结

  1. 环境检查:启动 Spark 集群和 HDFS,验证进程与权限;
  2. 数据准备:创建 HDFS 目录 → 准备本地文本 → 上传至 HDFS;
  3. 代码实现:按 Maven 规范创建目录(src/main/scala)→ 编写 Scala 代码;
  4. 打包验证:通过 Maven 编译打包,确保生成含主类的 JAR;
  5. 任务提交:上传 JAR 到 HDFS → 删除旧输出 → 执行 spark-submit
    4 输出目录已存在(FileAlreadyExistsException)

解决:提交前删除旧输出目录:hdfs dfs -rm -r /spark-wordcount/output

七、流程总结

  1. 环境检查:启动 Spark 集群和 HDFS,验证进程与权限;
  2. 数据准备:创建 HDFS 目录 → 准备本地文本 → 上传至 HDFS;
  3. 代码实现:按 Maven 规范创建目录(src/main/scala)→ 编写 Scala 代码;
  4. 打包验证:通过 Maven 编译打包,确保生成含主类的 JAR;
  5. 任务提交:上传 JAR 到 HDFS → 删除旧输出 → 执行 spark-submit
  6. 结果查看:通过 HDFS 命令或 Web UI 验证词频统计结果。

更多推荐