Spark 集群搭建与任务提交完整指南
·
🖥️ 环境信息
- 操作系统: macOS
- Java 版本: OpenJDK 11
- Scala 版本: 2.12.18
- Spark 版本: 3.5.0
- Hadoop 版本: 3.4.2
- 开发工具: IntelliJ IDEA
- 构建工具: Maven
🔧 遇到的问题及解决方案
问题 1: Spark 包结构识别问题
现象: com.example 文件夹没有被识别为 Scala 包
解决方案:
- 右键点击
scala文件夹 - 选择 "Mark Directory as" → "Sources Root"
- 或在 IDEA 中配置:File → Project Structure → Modules → Sources
问题 2: SparkContext 启动失败
错误信息:
SparkContext has been shutdown
原因: 设置了集群模式 spark://localhost:7077,但本地 Spark 集群未启动
解决方案: 启动 Spark 独立集群(见后续步骤)
问题 3: Spark Master 启动权限错误
错误信息:
javax.security.auth.Subject.getSubject
UserGroupInformation.getCurrentUser
原因: 使用 sudo 启动导致用户身份混乱
解决方案:
bash
# ❌ 错误做法
sudo ./sbin/start-master.sh
# ✅ 正确做法
# 1. 修改目录权限
sudo chown -R username:staff /opt/bigdata/spark-3.5.0
sudo chmod -R 755 /opt/bigdata/spark-3.5.0
# 2. 设置环境变量并启动(不用 sudo)
export HADOOP_USER_NAME=username
./sbin/start-master.sh
问题 4: Lambda 序列化错误(最关键的问题)
错误信息:
java.lang.ClassCastException: cannot assign instance of java.lang.invoke.SerializedLambda
to field org.apache.spark.rdd.MapPartitionsRDD.f of type scala.Function3
原因:
- IDEA 运行环境和 Spark 集群环境的 Java/Scala 版本序列化方式不兼容
- Lambda 表达式在跨 JVM 传输时反序列化失败
尝试过的方案:
- ❌ 添加 Kryo 序列化器 - 无效
- ❌ 统一 Java 版本 - 问题依旧
- ✅ 使用 spark-submit 打包提交 - 成功!
✅ 最终成功方案
核心思路
- 开发调试阶段: 在 IDEA 中使用
local[*]本地模式 - 集群运行阶段: 使用
spark-submit命令提交打包后的 JAR
📝 完整配置步骤
Step 1: 配置 Spark 环境
1.1 创建并配置 spark-env.sh
bash
cd /opt/bigdata/spark-3.5.0/conf
cp spark-env.sh.template spark-env.sh
vim spark-env.sh
添加以下内容:
bash
# 设置 Java 路径
export JAVA_HOME=/opt/homebrew/opt/openjdk@11/libexec/openjdk.jdk/Contents/Home
# 设置 Hadoop 用户
export HADOOP_USER_NAME=username
# Spark Master 配置
export SPARK_MASTER_HOST=localhost
export SPARK_MASTER_PORT=7077
export SPARK_MASTER_WEBUI_PORT=8080
# Worker 配置
export SPARK_WORKER_CORES=2
export SPARK_WORKER_MEMORY=2g
export SPARK_WORKER_WEBUI_PORT=8081
1.2 启动 Spark 集群
bash
cd /opt/bigdata/spark-3.5.0
# 设置环境变量
export HADOOP_USER_NAME=username
# 启动 Master
./sbin/start-master.sh
# 启动 Worker(使用实际的 Master URL)
./sbin/start-worker.sh spark://usernamedeMacBook-Air.local:7077
# 验证启动
jps | grep -E "Master|Worker"
1.3 访问 Spark Web UI
- Master UI: http://localhost:8080
- Worker UI: http://localhost:8081
Step 2: 配置 Maven 项目
2.1 pom.xml 配置
xml
<?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.example</groupId>
<artifactId>spark-demo</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<maven.compiler.source>11</maven.compiler.source>
<maven.compiler.target>11</maven.compiler.target>
<scala.version>2.12.18</scala.version>
<spark.version>3.5.0</spark.version>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<dependencies>
<!-- Scala 库 -->
<dependency>
<groupId>org.scala-lang</groupId>
<artifactId>scala-library</artifactId>
<version>${scala.version}</version>
</dependency>
<!-- Spark Core -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_2.12</artifactId>
<version>${spark.version}</version>
</dependency>
<!-- Spark SQL -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_2.12</artifactId>
<version>${spark.version}</version>
</dependency>
</dependencies>
<build>
<plugins>
<!-- Scala 编译插件 -->
<plugin>
<groupId>net.alchim31.maven</groupId>
<artifactId>scala-maven-plugin</artifactId>
<version>4.8.1</version>
<executions>
<execution>
<goals>
<goal>compile</goal>
<goal>testCompile</goal>
</goals>
</execution>
</executions>
</plugin>
<!-- 打包插件 -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.4.1</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>
<transformers>
<transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<mainClass>com.example.SparkDemo</mainClass>
</transformer>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
Step 3: 编写 Spark 应用代码
3.1 开发调试版本(IDEA 本地运行)
scala
package com.example
import org.apache.spark.{SparkConf, SparkContext}
object SparkDemo {
def main(args: Array[String]): Unit = {
val conf = new SparkConf()
.setAppName("Spark Demo")
.setMaster("local[*]") // 本地模式,用于开发调试
val sc = new SparkContext(conf)
try {
val data = 1 to 10000
val distData = sc.parallelize(data)
val result = distData.filter(_ < 10).collect()
println("过滤结果:")
result.foreach(println)
println(s"总共有 ${result.length} 个小于 10 的数字")
} finally {
sc.stop()
}
}
}
3.2 集群提交版本(spark-submit)
scala
package com.example
import org.apache.spark.{SparkConf, SparkContext}
object SparkDemo {
def main(args: Array[String]): Unit = {
val conf = new SparkConf()
.setAppName("Spark Demo")
// 不设置 master,由 spark-submit 参数指定
val sc = new SparkContext(conf)
try {
val data = 1 to 10000
val distData = sc.parallelize(data)
val result = distData.filter(_ < 10).collect()
println("过滤结果:")
result.foreach(println)
println(s"总共有 ${result.length} 个小于 10 的数字")
} finally {
sc.stop()
}
}
}
Step 4: 打包并提交任务
4.1 打包项目
bash
cd ~/spark-demo
mvn clean package
成功后生成:target/spark-demo-1.0-SNAPSHOT.jar
4.2 提交到 Spark 集群
bash
/opt/bigdata/spark-3.5.0/bin/spark-submit \
--class com.example.SparkDemo \
--master spark://usernamedeMacBook-Air.local:7077 \
--deploy-mode client \
--executor-memory 1g \
--total-executor-cores 2 \
~/spark-demo/target/spark-demo-1.0-SNAPSHOT.jar
参数说明:
--class: 主类名(包含包路径)--master: Spark Master 地址--deploy-mode:client: Driver 运行在本地(可以看到输出)cluster: Driver 运行在集群(生产环境)
--executor-memory: 每个 Executor 的内存--total-executor-cores: 使用的总核心数
🔍 常见问题排查
1. 查看 Spark 进程
bash
jps | grep -E "Master|Worker"
2. 查看 Master 日志
bash
tail -f /opt/bigdata/spark-3.5.0/logs/spark-*-Master-*.out
3. 查看 Worker 日志
bash
tail -f /opt/bigdata/spark-3.5.0/logs/spark-*-Worker-*.out
4. 停止 Spark 集群
bash
cd /opt/bigdata/spark-3.5.0
# 停止所有
./sbin/stop-all.sh
# 或单独停止
./sbin/stop-worker.sh
./sbin/stop-master.sh
5. 重启 Spark 集群
cd /opt/bigdata/spark-3.5.0
# 停止
./sbin/stop-all.sh
# 启动
export HADOOP_USER_NAME=username
./sbin/start-master.sh
./sbin/start-worker.sh spark://usernamedeMacBook-Air.local:7077
6. 检查端口占用
# 检查 7077 端口(Master)
lsof -i :7077
# 检查 8080 端口(Master Web UI)
lsof -i :8080
# 检查 8081 端口(Worker Web UI)
lsof -i :8081
📊 开发工作流程总结
┌─────────────────────────────────────────────────────────┐
│ 开发阶段 │
│ 1. 在 IDEA 中编写代码 │
│ 2. 使用 local[*] 模式本地调试 │
│ 3. 快速验证业务逻辑 │
└─────────────────────────────────────────────────────────┘
↓
┌─────────────────────────────────────────────────────────┐
│ 打包阶段 │
│ 1. 修改代码,移除或注释 setMaster │
│ 2. mvn clean package │
│ 3. 生成 JAR 文件 │
└─────────────────────────────────────────────────────────┘
↓
┌─────────────────────────────────────────────────────────┐
│ 提交阶段 │
│ 1. 确保 Spark 集群运行正常 │
│ 2. 使用 spark-submit 提交 JAR │
│ 3. 在 Web UI 监控任务执行 │
└─────────────────────────────────────────────────────────┘
💡 最佳实践建议
1. 本地开发模式
scala
// 适合:快速开发、单元测试、调试
.setMaster("local[*]")
2. 集群提交模式
bash
# 适合:生产环境、大数据处理、分布式计算
spark-submit --master spark://host:7077 your-app.jar
3. 版本一致性
确保以下组件版本一致:
- Scala 版本(2.12.x)
- Spark 版本(3.5.0)
- Java 版本(11)
4. 资源配置
根据实际情况调整:
bash
--executor-memory 2g # Executor 内存
--driver-memory 1g # Driver 内存
--total-executor-cores 4 # 总核心数
--executor-cores 2 # 每个 Executor 核心数
🎯 关键要点总结
- 不要用 sudo 启动 Spark - 会导致权限和用户身份问题
- Lambda 序列化问题 - 使用 spark-submit 提交 JAR 是最可靠的方案
- 环境变量很重要 - 启动前设置
HADOOP_USER_NAME - 两种模式分工明确:
local[*]: IDEA 开发调试spark://host:port: 集群生产运行
- 查看 Web UI - 实时监控任务执行情况和资源使用
📚 参考资源
更多推荐
所有评论(0)