
简介这是面向大数据课程实验的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 为什么实验要求里总带着HadoopSpark只是来借资源的几乎每份“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 21 | grep Using提示如果你在实验报告里写“启动Spark集群”老师会知道你还没搞懂——Spark通常不是“启动”的而是“提交作业到已有集群”。2.2 最小集群搭建每台机器必须改的三个配置假设你拿到的是三台虚拟机节点规划是master1核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_HOSTmaster export SPARK_WORKER_CORES2 export SPARK_WORKER_MEMORY3g export SPARK_DRIVER_MEMORY1g三个配置的用途要写进实验报告JAVA_HOME是Spark启动JVM的前提workers文件决定Master向哪些节点分发Executor进程SPARK_WORKER_MEMORY控制单台Worker能启动的Executor总内存这个参数要是设得比机器物理内存还大启动时不会报错作业一提交就开始OOM。紧接着验证集群状态。注意Spark的Web UI端口是8080Master和4040App临时端口很多人会把它和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总内存只有3gExecutor启动会卡在“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 JavaRDDString lines spark.read().textFile(hdfs://master:9000/user/hadoop/input/words.txt).javaRDD(); JavaRDDString words lines.flatMap(line - Arrays.asList(line.split( )).iterator()); JavaPairRDDString, Integer pairs words.mapToPair(word - new Tuple2(word, 1)); JavaPairRDDString, 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时才真正执行。// 这个写法打印不出任何东西 JavaRDDString upper lines.map(line - { System.out.println(map executing); return line.toUpperCase(); }); // 必须加行动算子才能看到打印 ListString 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减少分区不shufflerepartition增加分区shuffle JavaRDDString repartitioned lines.repartition(4); JavaRDDString 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对象没问题。但真实数据里如果某一行缺字段比如没有courseSpark推断出来的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如果你指定IntegerTypeSpark在转换时会先把脏值置为null默认mode是PERMISSIVE不会让作业直接崩溃——这在实际项目中更常用。提示如果你的实验允许用PythoncreateDataFrame接收的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默认是8192MB8G但你的虚拟机可能只有4G内存。Nodemanager发现自己可分配内存小于应用申请的内存应用就永远处于ACCEPTED状态。解决方法是把yarn-site.xml里的这个数调小property nameyarn.nodemanager.resource.memory-mb/name value3072/value /property5.3 现象端口被占用Spark Master起不来如果你先启动了Hadoop再启动Spark大概率没问题但如果先启动Spark然后Hadoop的SecondaryNameNode把50090端口占了Spark Master默认8080报了Address already in use你需要换端口。修改spark-env.sh里的SPARK_MASTER_PORT8088或直接关掉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撑爆。这三件事想清楚作业基本不会有大问题。这门实验的价值不在于那几十行代码而在于它逼你建立了“关注数据位置与数据分布”的意识。希望帮到你。本文还有配套的精品资源点击获取