Spark 3.3.0 词频统计全流程实操指南(Scala 版,适配 Hadoop 集群)
·
Spark 3.3.0 词频统计全流程实操指南(Scala 版,适配 Hadoop 集群)
一、环境与配置说明
1.1 集群环境参数
| 项目 | 配置详情 | 说明 |
|---|---|---|
| 集群拓扑 | Master:hadoop102(192.168.10.102);Worker:hadoop103、hadoop104 | 3 节点 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)
解决:
- 确保代码在
src/main/scala目录; - 验证
pom.xml中<mainClass>WordCount</mainClass>与代码类名一致; - 重新打包并覆盖 HDFS 上的 JAR。
6.3 结果缺失(如缺少 Hello、Spark)
解决:
- 验证 HDFS 输入文件内容:
hdfs dfs -cat /spark-wordcount/input/wordcount-data.txt; - 检查提交命令的输入路径是否正确(需指向
/spark-wordcount/input)。
6.4 输出目录已存在(FileAlreadyExistsException)
解决:提交前删除旧输出目录:hdfs dfs -rm -r /spark-wordcount/output。
七、流程总结
- 环境检查:启动 Spark 集群和 HDFS,验证进程与权限;
- 数据准备:创建 HDFS 目录 → 准备本地文本 → 上传至 HDFS;
- 代码实现:按 Maven 规范创建目录(
src/main/scala)→ 编写 Scala 代码; - 打包验证:通过 Maven 编译打包,确保生成含主类的 JAR;
- 任务提交:上传 JAR 到 HDFS → 删除旧输出 → 执行
spark-submit;
4 输出目录已存在(FileAlreadyExistsException)
解决:提交前删除旧输出目录:hdfs dfs -rm -r /spark-wordcount/output。
七、流程总结
- 环境检查:启动 Spark 集群和 HDFS,验证进程与权限;
- 数据准备:创建 HDFS 目录 → 准备本地文本 → 上传至 HDFS;
- 代码实现:按 Maven 规范创建目录(
src/main/scala)→ 编写 Scala 代码; - 打包验证:通过 Maven 编译打包,确保生成含主类的 JAR;
- 任务提交:上传 JAR 到 HDFS → 删除旧输出 → 执行
spark-submit; - 结果查看:通过 HDFS 命令或 Web UI 验证词频统计结果。
更多推荐
所有评论(0)