Spark初级编程实践:从spark-shell到独立应用的闭环之路 简介这是一份《大数据技术原理与应用》课程实验七的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.04Hadoop 3.1.3JDK 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_NAMEhadoop保证往 HDFS 写文件时用户身份一致export SPARK_LOCAL_IP127.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.txtSpark 解析时把它当成非法 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-shellsc.textFile(hdfs:///user/hadoop/test.txt).count()这里最容易错的是根目录概念。HDFS 的根是/user/hadoop/而不是~因为~是 shell 里的家目录展开符只对当前登录用户生效而 Spark 应用跑在 JVM 进程里不会替你展开~。报告里第二个坑正是这个写路径时用了~HDFS 里根本没有这个目录于是报InvalidInputException提示路径不存在。本地路径和 HDFS 路径的差异可以用一张表总结数据源推荐写法根路径文件不存在时报错Linux 本地文件file:///home/hadoop/test.txtLinux 根/FileNotFoundException或IllegalArgumentExceptionHDFS 文件hdfs:///user/hadoop/test.txtHDFS 根目录/InvalidInputExceptionHDFS 文件带 NameNode 地址hdfs://localhost:9000/user/hadoop/test.txtHDFS 根目录/同上如果 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.sbtname : 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.11spark-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-memoryDriver 内存数据量大时调大--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 第二个坑InvalidInputExceptionHDFS 根目录找错地方现象读 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 第三个坑URISyntaxExceptionURL 前面藏着看不见的空格现象写独立应用读取 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 里出错率奇高偏偏报错信息又很抽象与其记每个异常码不如把检查动作前置。希望你跑这套实验时能一次通过如果一不小心也栽在斜杠或空格上这份检查序列应该能帮你少走一段弯路。本文还有配套的精品资源点击获取