1. 环境准备:Ubuntu系统与基础工具

在开始构建Spark应用之前,我们需要确保Ubuntu系统已经准备好基础环境。我推荐使用Ubuntu 20.04 LTS或22.04 LTS版本,这两个版本长期支持且稳定性较好。如果你是全新安装的系统,建议先执行系统更新:

sudo apt update && sudo apt upgrade -y

接下来安装必要的开发工具链。这里有个小技巧:一次性安装所有常用工具可以避免后续频繁中断安装流程。我通常会执行以下命令:

sudo apt install -y openjdk-11-jdk curl gnupg2 software-properties-common

选择Java 11是因为目前大多数Spark版本都兼容JDK 11。安装完成后,验证Java版本:

java -version

如果看到类似"openjdk 11.0.xx"的输出,说明安装成功。这里有个容易踩的坑:有些系统可能预装了多个Java版本,可以通过update-alternatives --config java来切换默认版本。

注意:如果你使用云服务器,建议先配置SSH密钥登录并设置防火墙规则。我遇到过因为忘记开放端口导致后续Spark UI无法访问的情况。

2. 安装与配置sbt工具

2.1 官方方式安装sbt

sbt是Scala项目的标准构建工具,官方提供了多种安装方式。我最推荐使用deb包安装,这是最省心的方式:

echo "deb https://repo.scala-sbt.org/scalasbt/debian all main" | sudo tee /etc/apt/sources.list.d/sbt.list
echo "deb https://repo.scala-sbt.org/scalasbt/debian /" | sudo tee /etc/apt/sources.list.d/sbt_old.list
curl -sL "https://keyserver.ubuntu.com/pks/lookup?op=get&search=0x2EE0EA64E40A89B84B2DF73499E82A75642AC823" | sudo apt-key add
sudo apt update
sudo apt install -y sbt

安装完成后验证版本:

sbt sbtVersion

第一次运行会下载依赖,可能需要几分钟。这里有个实用技巧:可以预先配置镜像源加速下载。创建~/.sbt/repositories文件并添加以下内容:

[repositories]
local
maven-central: https://maven.aliyun.com/repository/central
typesafe: https://repo.typesafe.com/typesafe/ivy-releases/, [organization]/[module]/(scala_[scalaVersion]/)(sbt_[sbtVersion]/)[revision]/[type]s/[artifact](-[classifier]).[ext]

2.2 手动安装sbt(备用方案)

如果官方源安装失败,可以采用手动方式。我最近在阿里云ECS上就遇到过这种情况:

wget https://github.com/sbt/sbt/releases/download/v1.9.0/sbt-1.9.0.tgz
tar -zxvf sbt-1.9.0.tgz
sudo mv sbt /usr/local/

然后创建启动脚本/usr/local/bin/sbt

#!/bin/bash
SBT_OPTS="-Xms512M -Xmx1536M -Xss1M -XX:+CMSClassUnloadingEnabled"
java $SBT_OPTS -jar /usr/local/sbt/bin/sbt-launch.jar "$@"

别忘了给执行权限:

chmod +x /usr/local/bin/sbt

3. 创建第一个Spark项目

3.1 初始化项目结构

现在我们来创建标准的sbt项目结构。我习惯在~/projects目录下工作:

mkdir -p ~/projects/spark-demo/src/main/scala
cd ~/projects/spark-demo

创建build.sbt文件,这是项目的核心配置文件。新手常犯的错误是版本不兼容,这里给出经过验证的组合:

name := "Spark Demo"
version := "1.0"
scalaVersion := "2.12.15"
libraryDependencies += "org.apache.spark" %% "spark-core" % "3.3.2"

提示:Spark 3.x通常需要Scala 2.12,而Spark 2.x需要Scala 2.11。版本不匹配会导致各种奇怪的错误。

3.2 编写WordCount示例

src/main/scala目录下创建WordCount.scala:

import org.apache.spark.sql.SparkSession

object WordCount {
  def main(args: Array[String]) {
    val spark = SparkSession.builder
      .appName("Word Count")
      .master("local[*]")  // 使用所有可用核心
      .getOrCreate()
    
    val textFile = spark.sparkContext.textFile("README.md")  // 统计本项目README
    val counts = textFile.flatMap(line => line.split(" "))
                 .map(word => (word, 1))
                 .reduceByKey(_ + _)
    
    counts.collect().foreach(println)
    spark.stop()
  }
}

这个经典示例展示了Spark的核心操作:读取文本、分词、统计词频。我在教学时发现,初学者最容易困惑的是RDD的转换操作链式调用,建议逐步调试理解。

4. 构建与运行Spark应用

4.1 打包应用程序

在项目根目录执行打包命令:

sbt package

成功后会生成target/scala-2.12/spark-demo_2.12-1.0.jar。这里有个实用技巧:可以添加assembly插件生成包含所有依赖的fat jar。在project/plugins.sbt中添加:

addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.1.1")

然后在build.sbt中添加:

assembly / assemblyMergeStrategy := {
  case PathList("META-INF", xs @ _*) => MergeStrategy.discard
  case x => MergeStrategy.first
}

使用sbt assembly命令生成完整jar包。

4.2 提交Spark作业

假设你已经安装了Spark(如果没有,可以通过sudo apt install -y spark安装),运行:

spark-submit \
  --class "WordCount" \
  --master local[*] \
  target/scala-2.12/spark-demo_2.12-1.0.jar

如果想在集群运行,只需修改--master参数,比如:

spark-submit \
  --class "WordCount" \
  --master spark://your-master:7077 \
  --executor-memory 2G \
  --total-executor-cores 4 \
  target/scala-2.12/spark-demo_2.12-1.0.jar

我在生产环境中发现,资源参数需要根据集群实际情况调整。过大的内存设置反而会导致GC问题。

5. 进阶:处理结构化数据

让我们升级到Spark SQL处理CSV数据。首先添加依赖:

libraryDependencies += "org.apache.spark" %% "spark-sql" % "3.3.2"

创建新的示例程序EmployeeAnalysis.scala:

import org.apache.spark.sql.SparkSession

object EmployeeAnalysis {
  def main(args: Array[String]) {
    val spark = SparkSession.builder
      .appName("Employee Analysis")
      .getOrCreate()
    
    import spark.implicits._
    
    val df = spark.read
      .option("header", "true")
      .option("inferSchema", "true")
      .csv("employees.csv")  // 假设有这个文件
    
    df.createOrReplaceTempView("employees")
    
    // 执行SQL查询
    val results = spark.sql("""
      SELECT department, AVG(salary) as avg_salary 
      FROM employees 
      GROUP BY department
    """)
    
    results.show()
    spark.stop()
  }
}

这个示例展示了DataFrame API和SQL查询的结合使用。实际项目中,我建议明确指定schema而不是依赖类型推断,特别是在生产环境中。

6. 调试与优化技巧

6.1 查看Spark UI

Spark应用启动后,默认会在4040端口提供Web UI。我经常用它来:

  1. 查看执行计划可视化
  2. 检查各阶段耗时
  3. 识别数据倾斜问题

如果是本地模式,访问http://localhost:4040即可。有个小技巧:可以设置spark.ui.port参数指定端口:

.config("spark.ui.port", "4041")

6.2 日志配置

创建src/main/resources/log4j.properties文件控制日志级别:

log4j.rootCategory=WARN, console
log4j.appender.console=org.apache.log4j.ConsoleAppender
log4j.appender.console.target=System.err
log4j.appender.console.layout=org.apache.log4j.PatternLayout
log4j.appender.console.layout.ConversionPattern=%d{yy/MM/dd HH:mm:ss} %p %c{1}: %m%n

# 设置Spark相关日志级别
log4j.logger.org.apache.spark=WARN
log4j.logger.org.eclipse.jetty=WARN

6.3 性能优化建议

根据我的项目经验,以下几点能显著提升Spark性能:

  1. 合理分区:确保每个分区处理200MB左右数据
    df.repartition(100)  // 根据数据量调整
    
  2. 缓存重用:对频繁使用的DataFrame进行缓存
    df.cache()
    
  3. 广播变量:对小数据集使用广播
    val broadcastVar = spark.sparkContext.broadcast(lookupData)
    

7. 常见问题解决

问题1:sbt下载依赖极慢甚至失败

解决方案:除了前面提到的镜像配置,还可以设置环境变量:

export SBT_OPTS="-Dsbt.override.build.repos=true"

问题2:Spark作业报内存错误

尝试调整这些参数:

spark-submit \
  --conf spark.executor.memory=4G \
  --conf spark.driver.memory=2G \
  ...

问题3:类找不到(ClassNotFoundException)

确保打包时包含所有依赖,或者使用assembly插件生成fat jar。我曾经花了半天时间才发现是依赖冲突导致的这个问题。

问题4:版本兼容性问题

这是我整理的版本兼容表:

Spark版本Scala版本JDK版本
3.3.x2.128/11
3.0.x2.128/11
2.4.x2.118

最后提醒一点:Spark应用开发是个迭代过程,建议从小数据量开始测试,逐步扩展到全量数据。我在实际项目中通常会先采样1%的数据验证逻辑正确性

更多推荐