🖥️ 环境信息

  • 操作系统: macOS
  • Java 版本: OpenJDK 11
  • Scala 版本: 2.12.18
  • Spark 版本: 3.5.0
  • Hadoop 版本: 3.4.2
  • 开发工具: IntelliJ IDEA
  • 构建工具: Maven

🔧 遇到的问题及解决方案

问题 1: Spark 包结构识别问题

现象: com.example 文件夹没有被识别为 Scala 包

解决方案:

  1. 右键点击 scala 文件夹
  2. 选择 "Mark Directory as" → "Sources Root"
  3. 或在 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 传输时反序列化失败

尝试过的方案:

  1. ❌ 添加 Kryo 序列化器 - 无效
  2. ❌ 统一 Java 版本 - 问题依旧
  3. 使用 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

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 核心数

🎯 关键要点总结

  1. 不要用 sudo 启动 Spark - 会导致权限和用户身份问题
  2. Lambda 序列化问题 - 使用 spark-submit 提交 JAR 是最可靠的方案
  3. 环境变量很重要 - 启动前设置 HADOOP_USER_NAME
  4. 两种模式分工明确:
    • local[*]: IDEA 开发调试
    • spark://host:port: 集群生产运行
  5. 查看 Web UI - 实时监控任务执行情况和资源使用

📚 参考资源

更多推荐