从零到一:在Ubuntu上部署sbt并构建首个Spark应用
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。我经常用它来:
- 查看执行计划可视化
- 检查各阶段耗时
- 识别数据倾斜问题
如果是本地模式,访问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性能:
- 合理分区:确保每个分区处理200MB左右数据
df.repartition(100) // 根据数据量调整 - 缓存重用:对频繁使用的DataFrame进行缓存
df.cache() - 广播变量:对小数据集使用广播
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.x | 2.12 | 8/11 |
| 3.0.x | 2.12 | 8/11 |
| 2.4.x | 2.11 | 8 |
最后提醒一点:Spark应用开发是个迭代过程,建议从小数据量开始测试,逐步扩展到全量数据。我在实际项目中通常会先采样1%的数据验证逻辑正确性
更多推荐


所有评论(0)