简介:这是一份《大数据技术原理与应用》课程实验七的Spark初级编程实践完整实验报告,面向正在学习Spark与Hadoop的大数据初学者,也适合需要完成类似实验或复习Spark基础操作的高校学生。实验基于Windows 10宿主机与Ubuntu Kylin 16.04虚拟机环境,采用Hadoop 3.1.3、JDK 1.8,内容覆盖Spark安装与spark-shell启动,读取本地文件与HDFS文件统计行数,以及用Scala编写SimpleApp、RemDup、AvgScore三个独立应用,分别实现行数统计、两文件合并去重、多科目成绩平均值计算,并给出sbt打包与spark-submit提交的完整流程。压缩包内仅1个docx文档,大小1.9MB,包含实验环境说明、命令行操作截图、Scala代码与sbt配置,以及路径少写斜杠、HDFS路径错误、URL含空格等常见异常的解决办法。已有8338人学习浏览,适合需要快速上手Spark编程实践、排查环境与代码问题的读者参考。
1. Spark 入门最花时间的不是理论,是把这三个程序跑通
这份实验七 Spark 初级编程实践报告,把新手到能跑独立 Spark 应用的最后一公里讲得比较清楚:先让 spark-shell 能读本地文件和 HDFS,再用 Scala 写三个独立应用,最后用 sbt 打包、spark-submit 提交。我见过不少人在这一阶段卡住,卡点往往不在 RDD 概念上,而在路径字符串上——file:/// 少一个斜杠、HDFS 根目录认错、URL 前面多一个看不见的空格,三处小问题能各耗掉一个下午。这份报告把这三个报错连着解决方案都记下来了,又配了行数统计、去重、求平均成绩三个练手任务,适合刚装好 Spark、想找一份能照跑的完整闭环实验的人。
2. 先让 spark-shell 能读到数据:本地路径与 HDFS 路径的差异
2.1 环境怎么配:Hadoop 3.1.3 + JDK 1.8 + Ubuntu 16.04 的常见选型理由
实验环境写得很明确:宿主机 Windows 10 家庭版,虚拟机 Ubuntu Kylin 16.04,Hadoop 3.1.3,JDK 1.8,内存 16 GB。这套组合在课程实验里非常典型,原因是 Hadoop 3.1.3 和 Spark 2.x 对 JDK 8 的兼容性最稳,新版本 JDK 反而不一定省心。虚拟机跑的好处是方便做快照,配错了环境变量直接回滚,不用重装系统,这一点在入门阶段比性能更重要。
安装部分报告里只写了“解压至固定路径”,这是对的——Spark 是绿色安装,下载二进制包解压后设置环境变量就能跑。我一般会在/etc/profile或~/.bashrc里固定这几个变量:
export JAVA_HOME=/usr/lib/jvm/java-8-openjdk-amd64 export HADOOP_HOME=/usr/local/hadoop export SPARK_HOME=/usr/local/spark export PATH=$PATH:$JAVA_HOME/bin:$HADOOP_HOME/bin:$SPARK_HOME/bin参数说明:SPARK_HOME指向解压目录,PATH里加上spark-shell和spark-submit的所在路径,后续命令不用写全路径。环境变量改完执行source ~/.bashrc再验证。实际动手时我还会额外加两条:export HADOOP_USER_NAME=hadoop,保证往 HDFS 写文件时用户身份一致;export SPARK_LOCAL_IP=127.0.0.1,避免虚拟机多网卡时 Spark 找不到本机 IP。这两条在单机实验里不是必须,但能省掉不少奇怪的连接异常。
这份报告里启动命令是./bin/spark-shell,也就是在$SPARK_HOME目录下执行的相对路径写法。如果你喜欢任何位置都能启动,就用上面的PATH方式,直接敲spark-shell即可。
2.2 spark-shell 读本地文件:file:/// 的三斜杠到底缺了什么
第一个实操是读 Linux 本地文件/home/hadoop/test.txt并统计行数。在 spark-shell 里输入:
val textFile = sc.textFile("file:///home/hadoop/test.txt") textFile.count()第一行创建 RDD,第二行count()返回 Long 类型的行数。这里最值得记的不是count(),而是路径file:///home/hadoop/test.txt里那三个斜杠。file://是协议标识,/是空的主机名,第三个/才是 Linux 根路径的开始。所以file:///home/...合起来的意思是:本地文件系统、当前主机、/home/...这个绝对路径。
报告里踩的第一个坑就是少写一个斜杠,变成file://home/hadoop/test.txt,Spark 解析时把它当成非法 URI,直接抛IllegalArgumentException。这类异常日志很长,新手容易迷失在堆栈里,其实核心信息就是“路径格式不对”。判断方法很简单:在终端先用ls /home/hadoop/test.txt确认文件真实存在,再对照你代码里的字符串,三个斜杠一个都不能少。我个人的习惯是先在 spark-shell 里用sc.textFile(...).take(1)试读一行,确认能取出数据再写完整逻辑。
顺带提一个新手容易混淆的点:count()只返回行数,不返回内容。想看内容,用textFile.collect()或textFile.take(5)。刚接触 Spark 时总有人拿count()的结果去比对文件内容,发现对不上,其实是把两个算子搞混了。
2.3 spark-shell 读 HDFS 文件:/user/hadoop 不是 ~ 的等价路径
读 HDFS 比读本地文件多两步:先确认文件在不在,再写对路径。文件不存在时,先创建 HDFS 用户目录并上传:
hdfs dfs -mkdir -p /user/hadoop hdfs dfs -put /home/hadoop/test.txt /user/hadoop/test.txt hdfs dfs -ls /user/hadoop然后回到 spark-shell:
sc.textFile("hdfs:///user/hadoop/test.txt").count()这里最容易错的是根目录概念。HDFS 的根是/user/hadoop/而不是~,因为~是 shell 里的家目录展开符,只对当前登录用户生效;而 Spark 应用跑在 JVM 进程里,不会替你展开~。报告里第二个坑正是这个:写路径时用了~,HDFS 里根本没有这个目录,于是报InvalidInputException,提示路径不存在。
本地路径和 HDFS 路径的差异可以用一张表总结:
| 数据源 | 推荐写法 | 根路径 | 文件不存在时报错 |
|---|---|---|---|
| Linux 本地文件 | file:///home/hadoop/test.txt | Linux 根/ | FileNotFoundException或IllegalArgumentException |
| HDFS 文件 | hdfs:///user/hadoop/test.txt | HDFS 根目录/ | InvalidInputException |
| HDFS 文件(带 NameNode 地址) | hdfs://localhost:9000/user/hadoop/test.txt | HDFS 根目录/ | 同上 |
如果 NameNode 端口不是默认值,需要写完整地址hdfs://主机名:端口/user/hadoop/test.txt。单机实验里写hdfs:///user/hadoop/test.txt最省事,Spark 会从core-site.xml里自动读取fs.defaultFS。注意这里三个斜杠和本地文件的三斜杠含义不同:第一个是协议分隔,第二、三个斜杠组成了hdfs://这个 scheme 的固定格式。写 HDFS 路径时经常有人少写斜杠变成hdfs:/user/hadoop/...,同样会触发 URI 解析异常。
3. 三个 Scala 独立应用:SimpleApp、RemDup、AvgScore 的代码与提交流程
3.1 SimpleApp:读 HDFS 统计行数,sbt 打包的最小工程
spark-shell 适合验证想法,正式一点的做法是写独立应用程序,用 sbt 打包成 JAR,再用 spark-submit 提交。SimpleApp 的作用是读 HDFS 文件并统计行数,它展示了一个 Spark 应用的最小骨架。工程目录结构如下:
/usr/local/spark/mycode/HDFStest/ ├── src/main/scala/SimpleApp.scala ├── simple.sbt └── target/ # sbt package 后自动生成SimpleApp.scala 的完整内容:
import org.apache.spark.SparkContext import org.apache.spark.SparkConf object SimpleApp { def main(args: Array[String]): Unit = { val conf = new SparkConf().setAppName("SimpleApp") val sc = new SparkContext(conf) val textFile = sc.textFile("hdfs:///user/hadoop/test.txt") println("文件中行数为:" + textFile.count()) sc.stop() } }逻辑说明:SparkConf负责应用配置,setAppName只是给任务起名,不影响功能;SparkContext是 Spark 的入口,之后所有 RDD 操作都由它驱动;textFile读入文件后,count()返回 Long 行数并打印到 Driver 端控制台;最后sc.stop()释放资源。
再来是构建文件 simple.sbt:
name := "simple-hdfs-test" version := "1.0" scalaVersion := "2.11.8" libraryDependencies += "org.apache.spark" %% "spark-core" % "2.4.0"参数说明:name是工程名,它决定 JAR 包文件名前缀;scalaVersion要和本机 Spark 的 Scala 版本一致,报告里的 JAR 路径是target/scala-2.11/,说明用的是 Scala 2.11;spark-core的版本要和你的 Spark 安装版本匹配,不确定时在 spark-shell 启动日志里看 Spark version,写死一个不匹配的版本会在sbt package阶段报依赖解析失败。
打包和提交命令:
cd /usr/local/spark/mycode/HDFStest /usr/local/sbt/sbt package /usr/local/spark/bin/spark-submit --class "SimpleApp" \ /usr/local/spark/mycode/HDFStest/target/scala-2.11/a-simple-hdfs-test_2.11-1.0.jar打包完成后,JAR 文件默认落在target/scala-2.11/目录下,文件名由name字段自动生成,-会被转成_。提交时的--class参数要写主类名,如果代码里加了 package,比如com.example.SimpleApp,这里也要写全限定名com.example.SimpleApp,否则会报找不到主类。
3.2 RemDup:两个文件合并去重,distinct 背后的分区开销
第二个应用是数据去重:把文件 A 和 B 合并,剔除重复内容,输出到新文件 C。样例数据里每一行是“日期 字母”的组合,A、B 各自有重复日期但内容不同,合并后同一日期最多保留两条记录。
RemDup.scala 的典型实现:
import org.apache.spark.SparkContext import org.apache.spark.SparkConf object RemDup { def main(args: Array[String]): Unit = { val conf = new SparkConf().setAppName("RemDup") val sc = new SparkContext(conf) val fileA = sc.textFile("hdfs:///user/hadoop/A.txt") val fileB = sc.textFile("hdfs:///user/hadoop/B.txt") fileA.union(fileB).distinct().saveAsTextFile("hdfs:///user/hadoop/C") sc.stop() } }逻辑说明:union把两个 RDD 合并成一个,不做去重;distinct()做全局去重,这一步背后有 shuffle,框架会把相同内容通过网络汇总到同一分区再排重;saveAsTextFile把结果写到 HDFS 目录,注意它生成的是一个目录而非单个文件。
这里有个新手容易误解的点:distinct()的结果顺序和输入顺序不一定一致。样例输出 C 看起来是按日期排好的,但分布式计算里顺序本来就不保证,只要每一行的集合内容正确即可,不要拿“顺序对不对”来判断程序成败。判断去重结果对不对,应该统计总行数:C 的行数 = A 行数 + B 行数 - 两文件中完全相同的行数。如果想验证,用后面的getmerge把结果拉到本地再sort排序。
打包提交和 SimpleApp 流程一致:
cd /home/hadoop/sparkapp2/RemDup /usr/local/sbt/sbt package /usr/local/spark/bin/spark-submit --class "RemDup" \ /home/hadoop/sparkapp2/RemDup/target/scala-2.11/remove-duplication_2.11-1.0.jar如果数据量大,我一般会在union之前先对 A、B 各自做一次distinct(),再合并去重,减少 shuffle 的数据量。这个优化在入门阶段用不到,但养成“先减量再合并”的习惯是好的。
3.3 AvgScore:从多个成绩文件算平均分,解析一行里的多对记录
第三个应用是求学生平均成绩。输入有三个文件,每个文件是某学科的成绩,每行格式像Algorithm 成绩:小明 92 小红 87 小新 82 小丽 90,也就是说一行里有多组“名字 成绩”对。目标是对所有学生求三科平均分,输出到新文件。
AvgScore.scala 的实现:
import org.apache.spark.SparkContext import org.apache.spark.SparkConf object AvgScore { def main(args: Array[String]): Unit = { val conf = new SparkConf().setAppName("AvgScore") val sc = new SparkContext(conf) val inputPath = "hdfs:///user/hadoop/datas" val outputPath = "hdfs:///user/hadoop/avg_result" val rdd = sc.textFile(inputPath).flatMap { line => val scorePart = line.split("成绩:")(1) scorePart.trim.split("\\s+").grouped(2).map { pair => (pair(0), pair(1).toInt) }.toList } val sumCount = rdd.mapValues(score => (score, 1)) .reduceByKey((a, b) => (a._1 + b._1, a._2 + b._2)) val avg = sumCount.mapValues { case (sum, count) => BigDecimal(sum.toDouble / count).setScale(2, BigDecimal.RoundingMode.HALF_UP).toDouble } avg.saveAsTextFile(outputPath) avg.collect().foreach(println) sc.stop() } }逻辑说明:flatMap把每行拆成多组键值对,这里先按成绩:切分取出右侧内容,再用正则\\s+按空白切分,grouped(2)把字符串数组两两一组,变成(名字, 分数);mapValues把分数包装成(分数, 1),reduceByKey按名字聚合,得到总成绩和科目数;最后mapValues算平均值并用BigDecimal保留两位小数。
这段代码有两个要留意的边界。第一,如果某行格式不对、没有成绩:这个分隔符,line.split("成绩:")(1)会数组越界,稳妥写法是先filter(line => line.contains("成绩:"))再进入flatMap。第二,grouped(2)要求每组成绩严格成对出现,如果数据里有空行、多余空格,要先trim再处理。报告里给的样例数据恰好都是整齐的,但真实数据不会这么乖巧。
打包提交命令和前面两个一样,只是目录和主类名不同:
cd /home/hadoop/sparkapp3/AvgScore /usr/local/sbt/sbt package /usr/local/spark/bin/spark-submit --class "AvgScore" \ /home/hadoop/sparkapp3/AvgScore/target/scala-2.11/average-score_2.11-1.0.jar输出结果每个学生一行,格式是(名字,平均分)。注意报告的样例输出顺序是(小红,83.67)(小新,88.33)(小明,89.67)(小丽,88.67),和输入文件里名字出现的顺序不一致,这也是分布式聚合的常态,不要当成 bug 去排查。
3.4 spark-submit 提交的固定姿势与参数核对
三个应用跑下来,spark-submit的命令格式是固定的:
/usr/local/spark/bin/spark-submit [参数] <jar 包路径> [应用参数]常用参数如下表:
| 参数 | 作用 | 示例 |
|---|---|---|
--class | 指定主类名,带包名则写全限定名 | --class "SimpleApp" |
--master | 指定运行模式,不写默认 local | --master local[*]用满本地核 |
--driver-memory | Driver 内存,数据量大时调大 | --driver-memory 2g |
--executor-memory | 执行器内存,提交到集群时用 | --executor-memory 1g |
报告里的命令都写了--class "SimpleApp",双引号在 shell 里主要是防止类名被拆分,实际不带引号也能跑。单机实验不指定--master时,Spark 会落到 local 模式,够用。提交后如果应用秒退又没报错,先回看--class写没写对、JAR 路径是不是绝对路径。这三个应用踩的坑基本都在“主类名 + 路径”的组合上。
4. 常见问题与避坑:路径少个斜杠、多出空格、根目录认错
4.1 第一个坑:IllegalArgumentException,本地路径少了一个斜杠
现象:在 spark-shell 里执行sc.textFile("file://home/hadoop/test.txt"),立刻抛出IllegalArgumentException,日志里能看到 URI 解析相关的描述。
原因:file:///写成了file://,三斜杠少了一个。本地文件路径的正确 URI 格式是file:///绝对路径,前两个斜杠是file:协议的标准写法,第三个斜杠代表根目录。少一个斜杠后,home会被当成 authority(主机名),Spark 无法解析。
解决:把路径补成file:///home/hadoop/test.txt。核对方法:在终端先用ls /home/hadoop/test.txt确认文件位置,再对照代码里的字符串数一下file:后面到底有几个斜杠。这个坑特别隐蔽的原因是 IDE 或编辑器里看不出问题,只有运行到textFile才爆。
4.2 第二个坑:InvalidInputException,HDFS 根目录找错地方
现象:读 HDFS 文件时报InvalidInputException,异常信息里带着Input path does not exist,并且路径看起来不像预期目录。
原因:路径里写了~或者写成了/home/hadoop/test.txt。HDFS 的根目录是/,用户目录是/user/hadoop,~只在 shell 里会被展开成当前用户的家目录,Spark 的 JVM 进程不会做这个替换。换句话说,~在 spark-shell 里就是字面量~,HDFS 里自然找不到这个目录。
解决:先创建目录再上传文件,最后用hdfs dfs -ls /user/hadoop确认:
hdfs dfs -mkdir -p /user/hadoop hdfs dfs -put /home/hadoop/test.txt /user/hadoop/test.txt hdfs dfs -ls /user/hadoop确认无误后,代码里写hdfs:///user/hadoop/test.txt。如果是集群环境,把hdfs:///换成hdfs://namenode:port/。我每次写 HDFS 路径前都会先在命令行跑一遍hdfs dfs -ls,因为路径是否存在这种事,用眼睛看永远比猜靠谱。
4.3 第三个坑:URISyntaxException,URL 前面藏着看不见的空格
现象:写独立应用读取 HDFS 文件时,运行报java.net.URISyntaxException: Illegal character in scheme name at index 0。
原因:代码里 HDFS 地址字符串前面多了一个空格,比如" hdfs:///user/hadoop/test.txt"。Spark 解析 URI 时,第一个字符不是字母而是空格,scheme 名称里出现非法字符,直接抛异常。这个问题在代码里肉眼极难发现,因为空格和正常字符串在大多数编辑器里几乎没有视觉差异。
解决:把字符串首尾的空格删掉,或者统一加trim():
val path = " hdfs:///user/hadoop/test.txt".trim sc.textFile(path)这个坑给我留下的印象最深,因为报错信息里的index 0很容易让人误以为是下标越界,排查方向完全跑偏。后来我养成的习惯是:凡是路径字符串,先println出来看首尾有没有空格,再进textFile。
4.4 实战里还会遇到的几个资源类报错
除了报告里记录的三个路径坑,Spark 入门常见的还有这三类,现象和原因比较典型。
OOM(内存溢出):现象是应用跑到一半报java.lang.OutOfMemoryError,常见诱因是collect()把全量数据拉回 Driver,或者数据量超过默认驱动内存。解决方式是先take()抽样确认数据,再考虑增大内存:spark-submit --driver-memory 2g。单机实验配置不高时,不要轻易对大数据集collect()。
找不到主类:现象是提交后立刻报ClassNotFoundException,原因通常是--class里写的类名没有带包名,或者代码里改了包名但提交命令没同步。解决方式是在 sbt 工程里统一维护主类名,提交前看一下打包出的 JAR 里实际的类路径。
输出目录已存在:现象是saveAsTextFile报FileAlreadyExistsException。Spark 的输出 API 不允许覆盖已存在的目录,解决方式是换一个新目录名,或者先手动删除旧目录:
hdfs dfs -rm -r /user/hadoop/C这三个问题报告里没写,但属于同一阶段必然遇到的环境类报错,提前知道能省不少排查时间。
5. 验证结果的一个习惯:先手工算一遍,再让 Spark 给你交答案
5.1 用 wc -l、sort、uniq 把手算结果和 Spark 结果对一遍
Spark 跑出来不一定对,尤其是去重和平均值这类需要 shuffle 的操作。我的验证习惯是先在 shell 里用传统命令拿到精确答案,再拿 Spark 输出做比对,两边一致才认为程序正确。
读文件统计行数,可以用wc -l对照:
wc -l /home/hadoop/test.txt这个数字应该和sc.textFile(...).count()完全一致。不一致时优先查文件编码和换行符,比如文件最后一行没有换行符时,textFile的统计口径可能和wc -l差一行。
去重结果的验证,把 HDFS 上的输出合并成一个本地文件,再用sort和uniq -c统计:
hdfs dfs -getmerge /user/hadoop/C ./C_local.txt sort C_local.txt | uniq -cuniq -c会打印每行出现的次数,去重正确的输出里每行次数都应该是 1。同时核对总行数:A 行数 + B 行数 - 两边重复的对数。
平均值的验证更直接,把样例数据当作笔算题:小明三科成绩 92、95、82,平均 89.67;小红 87、81、83,平均 83.67;小新 82、89、94,平均 88.33;小丽 90、85、91,平均 88.67。Spark 输出保留两位小数,结果和手算一致才算通过。保留精度这一环我习惯放在程序里做,而不是输出后再格式化,避免不同计算环境下浮点尾数不一致。
5.2 一个值得固定下来的路径检查序列
这几轮跑下来,最大的收获不是学会了三个算子,而是建立了一套固定的路径检查序列。从那以后我每次写 Spark 应用,都会强制走一遍以下五步。
第一步,打印路径字符串,确认首尾没有空格。第二步,数file:///的斜杠数,本地路径必须是三个斜杠。第三步,HDFS 路径先执行hdfs dfs -ls,确认目录存在再写进代码。第四步,在 spark-shell 里用sc.textFile(...).take(1)验证,而不是直接丢进独立应用。第五步,提交前核对--class的类名和 JAR 包内实际主类一致。
这套序列听起来笨,但它把定位问题的范围从“日志堆栈 + 猜测”压缩到“路径字符串本身”。报告里记录的三个报错,全都在前四步里能暴露出来。路径这种东西,在 Spark 里出错率奇高,偏偏报错信息又很抽象,与其记每个异常码,不如把检查动作前置。希望你跑这套实验时能一次通过,如果一不小心也栽在斜杠或空格上,这份检查序列应该能帮你少走一段弯路。
本文还有配套的精品资源,点击获取