简介:这是面向大数据课程实验的Spark初级编程实践资源,适合正在学习Hadoop与Spark的学生或自学者参考。资源以单份docx实验报告形式呈现,压缩包仅1个文件,大小约1.9MB,内容完整记录了“大数据技术原理与应用”课程实验七的全部过程。报告从实验环境配置讲起,包括Ubuntu虚拟机下的Hadoop 3.1.3与JDK 1.8环境,再到Spark shell启动、本地及HDFS文件读取与行数统计,并给出Scala独立应用SimpleApp、RemDup和AvgScore的完整代码与打包运行命令,覆盖数据去重、平均值计算等典型场景。值得一提的是,报告还整理了三个常见报错(如路径缺斜杠、HDFS根目录识别错误、URL含空格)及对应解决方案,能帮助读者避开同类坑点。已有8338人学习下载,对于需要提交实验报告或快速梳理Spark入门操作的同学具有实用参考价值。
1. 实验七:Spark初级编程实践——这门课到底让你学会什么
如果你正在为“实验七:Spark初级编程实践”这门课发愁,大概率是卡在了“代码能跑但不知道为什么能跑”的阶段。这个实验在多数高校大数据课程里的定位很明确:它不是让你调参调优,也不是让你做复杂的数据管道,而是让你把Spark当成一个“分布式计算器”,亲手在上面跑通RDD转换、Action算子、以及一个完整的Spark SQL作业。做完它,你应该能回答三个问题:Spark和Hadoop到底什么关系、RDD为什么比普通数组难用、以及一个Spark作业从提交到出结果经历了什么。
这套实验通常放在Hadoop生态课程的中段,前置要求是Linux基本操作和Java/Scala语法,不需要你懂源码。但恰恰因为“看起来简单”,很多人会翻车在环境变量、端口占用、JSON解析这类基础问题上。这篇笔记按“环境搭建 → RDD编程 → Spark SQL实践 → 排错避坑 → 验证技巧”的顺序推进,每一条都是能直接抄作业的步骤和参数。
2. 集群环境搭建:先搞清楚Spark是“客人”不是“主人”
2.1 为什么实验要求里总带着Hadoop:Spark只是来借资源的
几乎每份“Spark初级编程实践”实验指导书开头都会写“请先确保Hadoop集群可用”。不少初学者在这里困惑:Spark不是自己的集群吗,为什么非要先装Hadoop?这个问题的答案,直接决定了你后面排错的方向。
Spark本身不承担分布式存储职责。它的计算模型可以跑在内存里,但数据来源和最终落地大多要依赖HDFS。实验场景里最常见的组合是:HDFS负责存实验数据,YARN负责分配CPU和内存资源,Spark作为计算框架向YARN申请资源并干活。也就是说,Spark是YARN的“租客”,不是房东。
在实验环境里,如果你用的是伪分布式Hadoop,那Spark也跑在伪分布式模式下,所有进程都堆在一台机器上。而如果实验要求“三台虚拟机搭建Spark集群”,那往往意味着NameNode和ResourceManager各管一摊,Worker节点上的Spark进程和DataNode进程共存。判断你的Spark作业到底提交给了谁,可以用下面命令看日志开头:
# 查看Spark作业运行模式,结果可能是yarn、standalone、local[*] spark-submit --version 2>&1 | grep "Using"提示:如果你在实验报告里写“启动Spark集群”,老师会知道你还没搞懂——Spark通常不是“启动”的,而是“提交作业到已有集群”。
2.2 最小集群搭建:每台机器必须改的三个配置
假设你拿到的是三台虚拟机,节点规划是master(1核2G)、worker1、worker2(各2核4G)。在下载好spark-3.3.x-bin-hadoop3这个版本之后,解压到 /opt/spark 下,然后必须要做三件事:配JAVA_HOME、配workers文件、配spark-env.sh。
# 1. 在 /etc/profile 里追加以下内容(三台机器都要做) export JAVA_HOME=/usr/lib/jvm/java-8-openjdk-amd64 export SPARK_HOME=/opt/spark export PATH=$PATH:$SPARK_HOME/bin:$SPARK_HOME/sbin # 2. 修改 $SPARK_HOME/conf/workers(原slaves文件在3.0后改名了) # 内容不要写 localhost,写两台worker的主机名 worker1 worker2 # 3. 修改 $SPARK_HOME/conf/spark-env.sh export SPARK_MASTER_HOST=master export SPARK_WORKER_CORES=2 export SPARK_WORKER_MEMORY=3g export SPARK_DRIVER_MEMORY=1g三个配置的用途要写进实验报告:JAVA_HOME是Spark启动JVM的前提;workers文件决定Master向哪些节点分发Executor进程;SPARK_WORKER_MEMORY控制单台Worker能启动的Executor总内存,这个参数要是设得比机器物理内存还大,启动时不会报错,作业一提交就开始OOM。
紧接着验证集群状态。注意,Spark的Web UI端口是8080(Master)和4040(App临时端口),很多人会把它和Hadoop的50070或9870搞混。
# 先启动HDFS,再启动Spark(顺序别反) start-dfs.sh $SPARK_HOME/sbin/start-master.sh $SPARK_HOME/sbin/start-workers.sh # 验证进程 jps # master上应有 Master 进程,worker1/worker2上应有 Worker 进程2.3 内存参数的分配逻辑:别把实验机器跑死
实验指导书很少告诉你该给Spark分多少内存,但它恰恰是作业能不能跑完的关键。一个通用经验是:给YARN和Spark的总内存之和,不要超过机器物理内存的75%。比如worker机器是4G内存,Hadoop的DataNode默认占用1G,那么Spark的SPARK_WORKER_MEMORY最好设在2g左右,留1G给操作系统和其他进程。
另外还有一个初学者几乎必踩的坑:在spark-submit时同时指定--num-executors 和 --executor-memory,会把所有资源挤爆。比如--executor-memory 2g但Worker总内存只有3g,Executor启动会卡在“Waiting for resources”状态,日志里反复出现“No sufficient resources”。这时候应该做的是:要么减少executor内存,要么增加Worker内存。
(我在2.2里的配置给它2g,是因为后面要跑Spark SQL读JSON,数据量虽然小但堆内存开销偏高,1g跑起来gc太频繁,日志里全是Full GC,作业慢得像是集群挂了。)
3. RDD编程实践:从WordCount到掌握算子的数据流向
3.1 用sc.textFile读数据的路径陷阱
实验的第一个编程题通常是WordCount,看着简单却最能暴露问题。先贴一份完整可跑的Java版本(实际上大多数实验允许用Python,但Java是入门大数据的最正宗姿势):
import org.apache.spark.api.java.JavaRDD; import org.apache.spark.api.java.JavaPairRDD; import org.apache.spark.sql.SparkSession; import scala.Tuple2; import java.util.Arrays; public class WordCount { public static void main(String[] args) { SparkSession spark = SparkSession.builder() .appName("WordCountLab") .master(args.length > 0 ? args[0] : "local[*]") .getOrCreate(); // 读HDFS上/user/hadoop/input/words.txt JavaRDD<String> lines = spark.read().textFile("hdfs://master:9000/user/hadoop/input/words.txt").javaRDD(); JavaRDD<String> words = lines.flatMap(line -> Arrays.asList(line.split(" ")).iterator()); JavaPairRDD<String, Integer> pairs = words.mapToPair(word -> new Tuple2<>(word, 1)); JavaPairRDD<String, Integer> counts = pairs.reduceByKey(Integer::sum); // 必须触发行动算子,前面全是惰性转换 counts.saveAsTextFile("hdfs://master:9000/user/hadoop/output/wc_result_" + System.currentTimeMillis()); spark.stop(); } }这段代码里有两个关键节点。第一个是textFile的路径:实验指导书如果是基于Hadoop写的,让你把文件放到HDFS的/input下,你没做这一步就改读本地路径file:///home/hadoop/words.txt,在伪分布式上可以,但在三节点集群上,Executor不在本地文件所在的那台机器时就会报FileNotFoundException——这个异常信息会让你误以为代码写错了,实际上是文件根本不在那个机器的本地磁盘上。
第二个是saveAsTextFile的输出目录必须不存在。Spark不会像普通Java程序那样帮你自动创建或覆盖目录,目录已存在会直接抛org.apache.hadoop.mapred.FileAlreadyExistsException。很多人的解决方案是手动删目录,更聪明的做法是像上面代码一样在输出路径后面拼时间戳,一劳永逸。
3.2 惰性求值和行动算子的关系:为什么你debug看不到效果
刚接触RDD的人普遍有一个困惑:我在map里加了System.out.println,为什么运行时不打印?原因是转换算子(map、flatMap、filter、reduceByKey)全是惰性求值的,Spark把它们当成一张“菜谱”,只有出现行动算子(collect、count、saveAsTextFile、foreach)时才真正执行。
// 这个写法打印不出任何东西 JavaRDD<String> upper = lines.map(line -> { System.out.println("map executing"); return line.toUpperCase(); }); // 必须加行动算子,才能看到打印 List<String> result = upper.collect();这里有一个实操技巧:在实验里如果你只想看前几条数据,不要用collect()把整个RDD拉回Driver。比如一个1GB的文件,collect()会试图把所有数据塞进Driver内存,直接OOM。标准做法是用take(10)看前10条,或者用foreachPartition在Executor端打印。
// 调试首选:take只拉少量数据到Driver upper.take(10).forEach(System.out::println);3.3 分区数对作业性能的影响:三个必调参数
实验数据量小的时候,分区数的影响完全看不出来,但你要是不理解,后面做网约车数据清洗这类真实项目时会很痛苦。RDD分区数主要由两个因素决定:输入文件切分大小,以及repartition/coalesce的显式调用。
// 强制重分区,coalesce减少分区(不shuffle),repartition增加分区(shuffle) JavaRDD<String> repartitioned = lines.repartition(4); JavaRDD<String> coalesced = lines.coalesce(2); // 查看当前分区数 System.out.println("Partitions: " + lines.getNumPartitions());在实验报告里,你应该能回答:textFile读HDFS文件时,默认块大小是128MB,也就是说一个120MB的文件只会有一个分区,整个Spark作业只用1个核在跑,其他worker都在围观。要利用集群的并行能力,需要repartition(3)或者调大spark.sql.files.maxPartitionBytes参数。(真实的网约车数据清洗项目那节课,我带了 3 个几十 GB 的日志文件,默认分区直接把单个 Executor 内存干爆,调大分区数之后问题立刻消失。)
注意:不要对每个RDD都盲目repartition,洗牌代价很高。你在Kafka分区数为6时,最好不要把Spark分区也硬设成和它一致,除非确实有下游并发需求。
4. Spark SQL与JSON读取:初级实验里最接近生产的环节
4.1 用DataFrame API读取JSON的三个坑
“spark中读取json”几乎是这门实验的标配。因为JSON是互联网数据的主流格式,课程设计者希望你能从结构化数据走向半结构化数据。先给一份Spark SQL完整案例——读取一个学生成绩JSON文件,算出平均分并按分数降序排:
// 输入文件 /user/hadoop/input/students.json {"name": "zhangsan", "score": 85, "course": "math"} {"name": "lisi", "score": 92, "course": "math"} {"name": "wangwu", "score": 78, "course": "python"} {"name": "zhaoliu", "score": 88, "course": "python"}# 用pyspark写更贴合初级实验场景 from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("JsonReadLab") \ .master("local[*]") \ .getOrCreate() df = spark.read.json("hdfs://master:9000/user/hadoop/input/students.json") df.printSchema() df.show() # 注册临时视图,用纯SQL做聚合——这是实验报告里加分的写法 df.createOrReplaceTempView("students") result = spark.sql(""" SELECT course, AVG(score) as avg_score, COUNT(1) as cnt FROM students GROUP BY course ORDER BY avg_score DESC """) result.show() # 写回成Parquet格式(实验7一般不做,提一下即可) result.write.mode("overwrite").parquet("hdfs://master:9000/user/hadoop/output/avg_score")这个实验正确的写法是用spark.read.json而不是spark.read.text再去手动解析,学历历上有三个坑。
第一个坑是schema推断。上面这个例子是每行一个完整JSON对象,没问题。但真实数据里如果某一行缺字段,比如没有course,Spark推断出来的schema会把course列变成nullable,聚合结果里会出现一行null值。处理方式是读数据时显式指定schema,或者用dropna(subset=["course"])。
第二个坑是多行JSON文件。有些文件是把一个JSON数组放在一行里:
[{"name":"zhangsan","score":85}, {"name":"lisi","score":92}]这种文件spark.read.json会解析失败,要么用spark.read.text读取后from_json解析,要么保持一致,把数据整理成JSON Lines格式。(JSON Lines是每行一个独立JSON对象,大数据领域标准格式,不需要引入额外的multiLine参数。)
第三个坑是压缩格式。JSON文件是.gz或.bz2结尾时,spark.read.json能自动识别不需要额外改代码;但如果是你自己用textFile读的,取回来的是压缩文件内部的文本流,按行拆分时容易把半行截断。
4.2 从RDD到DataFrame的转换方式对比
实验里经常会出现“给你一个RDD,要求转成DataFrame用SQL分析”的题目,需求场景是已有文本数据但想用SQL方式查询。三种主流做法各适用不同场景:反射推断(case class)、编程式指定schema、toDF加列名。
// Scala版本:反射推断(需要定义case class,适合列名和数据类型已知的情况) import org.apache.spark.sql.Encoders case class Student(name: String, score: Int, course: String) val rdd = spark.sparkContext.textFile("hdfs://master:9000/user/hadoop/input/students.csv") .map(_.split(",")) .map(p => Student(p(0), p(1).trim.toInt, p(2))) import spark.implicits._ val df = rdd.toDF() df.createOrReplaceTempView("students") val result = spark.sql("SELECT course, AVG(score) FROM students GROUP BY course") result.show()# Python版本:编程式schema——更灵活,适合字段名或类型在运行时才确定的场景 from pyspark.sql.types import StructType, StructField, StringType, IntegerType schema = StructType([ StructField("name", StringType(), True), StructField("score", IntegerType(), True), StructField("course", StringType(), True) ]) rdd = spark.sparkContext.textFile("hdfs://master:9000/user/hadoop/input/students.csv") \ .map(lambda line: line.split(",")) \ .map(lambda p: (p[0], int(p[1]), p[2])) df = spark.createDataFrame(rdd, schema) df.show()这两种方式的差异很微妙。反射推断只要case class里字段顺序和数据一致,写起来代码最短;但一旦数据里混进脏值,比如score字段出现了"85分"这种字符串,反射推断会在运行时粗暴地抛异常。编程式schema如果你指定IntegerType,Spark在转换时会先把脏值置为null(默认mode是PERMISSIVE),不会让作业直接崩溃——这在实际项目中更常用。
提示:如果你的实验允许用Python,
createDataFrame接收的RDD元素类型是tuple或list,不要传dict,否则会报ValueError: Unexpected tuple之类的误导性错误。
4.3 Spark SQL与Hive的边界:别把实验做成数仓
有些实验指导书会在“实验七”的进阶部分让你把Spark SQL结果存到Hive表里,这时候需要在spark-env.sh配HIVE_HOME,以及把hive-site.xml放到Spark的conf目录下。但你要清楚,这已经属于“数仓才需要的能力”而不是“初级编程实践”的目标。如果实验没有明确要求,不建议浪费时间在这个方向,因为版本兼容问题实在太多(Spark 3.3内置的Hive版本是2.3.9,和你的Hadoop集群Hive版本一旦不一致,spark.sql.warehouse.dir找不到表时会直接报Table not found)。
(说白了,读json文件的实验重点在schema推断、类型转换和视图注册,这三样做完,实验核心能力就已经具备了。)
5. 常见问题与排错:实验七翻车现象前五名,附解决办法
5.1 现象:日志里大量“Lost task”但作业最后还是成功
罪魁祸首通常是Executor内存不够,某个task处理数据时导致GC停顿,Spark的推测执行机制(speculation)自动在另一个节点重启了task。小数据量下感知不到,但日志里会有明显痕迹——“Lost task 0.0 in stage 0.0 (TID 2) on executor 1: ExecutorLostFailure”。
排查方式:先看是否是数据倾斜。用getNumPartitions看分区数,如果有任务处理了99%的数据,另一个任务处理1%,就是键分布不均。解决思路是加盐或改分区策略,但在初级实验里,最简单的处理是调大spark.executor.memory,或者用repartition(分区数调大)把数据切更碎。我的习惯是先看一眼是不是某个文件块异常大,这种场景直接用textFile的minPartitions参数解决:
// 强制至少拆成6个分区,避免单task数据过大 val rdd = spark.sparkContext.textFile("hdfs://.../big_file.txt", 6)5.2 现象:spark-submit 提交Python脚本卡在“YARN Application is in ACCEPTED state”
最常见的不是资源不足,而是你根本没给YARN分配足够资源。
在yarn-site.xml里,yarn.nodemanager.resource.memory-mb默认是8192MB(8G),但你的虚拟机可能只有4G内存。Nodemanager发现自己可分配内存小于应用申请的内存,应用就永远处于ACCEPTED状态。解决方法是把yarn-site.xml里的这个数调小:
<property> <name>yarn.nodemanager.resource.memory-mb</name> <value>3072</value> </property>5.3 现象:端口被占用,Spark Master起不来
如果你先启动了Hadoop,再启动Spark,大概率没问题;但如果先启动Spark,然后Hadoop的SecondaryNameNode把50090端口占了,Spark Master(默认8080)报了Address already in use,你需要换端口。
修改spark-env.sh里的SPARK_MASTER_PORT=8088,或直接关掉Hadoop的某个进程再重排启动顺序。从实验“正确性”角度说,推荐后者——让Spark使用默认端口,不给自己留不必要的变量。
5.4 现象:java.lang.OutOfMemoryError: Java heap space发生在Driver端
数据量明明很小,为什么Driver会OOM?原因是你在Driver端不小心把大RDD给collect()了。这不是调大SPARK_DRIVER_MEMORY能根治的——它就是1GB的数据,你怎么调都会爆。正确做法是改用saveAsTextFile或foreachPartition把数据落到外部。请时刻记住:Driver不是用来装数据的,是用来调度和写SQL的。
5.5 现象:读CSV文件时字符串列出现双引号未去除
原因:用spark.read.csv而不指定quote参数时,Spark默认认为双引号是转义符。当数据本身包含带引号的字段时,option("quote", "'")或者option("escape", "\"")能解决。这个在你的实验数据里不一定会出现,但如果你去读网约车数据,一列地址里夹杂几个双引号很正常。
(整套实验跑完,五个坑基本覆盖了至少80%的新手场上遇到的报错。)
6. 最后的实用技巧:用RDD转换验证实验结果的正确性
实验提交前,你总会怀疑结果对不对。有一个不依赖老师的自检方式:用Spark自身的两种API算同一份数据,比对结果。比如用RDD方式算WordCount,再用Spark SQL的GROUP BY算同一条数据,两边结果不一致,那说明你的逻辑在某个算子上有偏差。
# 比对两份结果(推荐用diff,而不是肉眼对比几十条记录) hdfs dfs -cat /user/hadoop/output/wc_result_*/part-* > /tmp/rdd_result.txt spark-sql --master yarn --executor-memory 1g \ -e "SELECT word, COUNT(1) FROM (SELECT explode(split(value, ' ')) AS word FROM textfile '/user/hadoop/input/words.txt') t GROUP BY word;" \ > /tmp/sql_result.txt 2>/dev/null diff /tmp/rdd_result.txt /tmp/sql_result.txt这个技巧的深层价值在于它逼你理解了Spark的两种API其实殊途同归——全都在底层转化为一组Task执行计划。实验报告里如果能写出“用两种方式交叉验证了数据正确性”,会比只贴一份输出结果要扎实得多。
再补一个简单的数据过滤验证技巧:如果你的Spark SQL结果里有大量null值,不要急着删。先用groupBy看有多少个null,再回源数据确认是源头脏数据还是你的关联键写得不对,这个习惯对后续接触更复杂的数据分析案例很有用。
落到习惯层面,我如今每次提交Spark作业前都会先想一遍三个问题:输入文件在哪个节点上,分区数大概多少,Driver内存会不会被collect撑爆。这三件事想清楚,作业基本不会有大问题。这门实验的价值不在于那几十行代码,而在于它逼你建立了“关注数据位置与数据分布”的意识。
希望帮到你。
本文还有配套的精品资源,点击获取