简介:本资源是《大数据技术原理与应用》课程配套的「Spark初级编程实践」实验报告,面向高校大数据初学者及Hadoop/Spark入门学习者,聚焦分布式计算框架的核心操作训练。内容覆盖Spark环境搭建(Ubuntu虚拟机+Hadoop 3.1.3+Spark+JDK 1.8)、Spark Shell交互式数据读取(本地文件与HDFS双路径)、独立Scala应用程序开发(SimpleApp行数统计、RemDup数据去重、AvgScore多文件平均分计算),并附详细排错指南(如file:///路径缺失、HDFS根目录误用、URL空格导致URISyntaxException等典型问题)。资源为1个1.9MB的DOCX文档,含完整实验环境配置、4类实验的操作步骤、代码截图、运行结果图示及问题分析,结构清晰、图文并茂,便于对照复现与理解原理。目前已有8338人学习下载,是Spark入门阶段不可多得的实操型教学参考材料。
1. Spark初级编程实践:不是配环境、跑WordCount就叫“会Spark”,而是搞懂Driver怎么发任务、Executor怎么扛住Shuffle、为什么本地模式跑通了集群却OOM
很多人把“Spark初级编程实践”当成一门配置课——装好Scala、下载Spark包、改几行spark-submit参数,跑出一个Hello World式的WordCount,就以为跨进了大数据门。但真实业务里,你刚把清洗脚本从本地spark-shell挪到YARN集群,任务就卡在Stage 0;你调大--executor-memory,集群反而报Container is running beyond physical memory limits;你用df.join()关联两张表,数据量一过千万,GC时间飙升到每分钟30秒……这些不是玄学,是Spark执行模型没吃透的必然翻车。本实践不讲“Spark是什么”,直奔一线工程师写第一个可交付作业必须踩实的三块地:本地伪分布式调试链路怎么闭环、RDD与DataFrame API选型边界在哪、Shuffle阶段内存和磁盘如何协同不拖垮任务。适合刚学完Scala基础、手上有台8G内存笔记本、想用真实小数据集(比如农产品价格CSV)练出肌肉记忆的开发者。别怕报错——每个报错背后都藏着Spark调度器的一次真实决策。
2. 用Spark 3.5.0在本地跑通WordCount:从解压到验证输出,一条命令闭环调试链
Spark初级实践的第一道门槛,从来不是语法,而是环境可信度。你不确定是代码写错了,还是SPARK_HOME指向了旧版本,或是JAVA_HOME用了JDK17而Spark只认JDK8——这种模糊地带会让新手在“Hello World”上卡三天。下面这套流程,是我给新人搭本地沙箱的标准动作:不依赖IDE插件、不碰任何配置文件、纯命令行验证每一步输出,确保你看到的count: 1247,就是Spark真正在算,不是缓存或mock。
2.1 下载、解压、验证Java与Spark版本兼容性
Spark 3.5.0官方明确要求JDK8或JDK11(注意:JDK17+需手动编译源码,生产环境不推荐)。先确认你的Java版本:
java -version # 必须输出类似:openjdk version "11.0.22" 2024-01-16 # 如果是17+,立刻切回JDK11:export JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64然后下载二进制包(不要用apt install spark——Ubuntu源里的Spark版本老旧且缺spark-sql模块):
wget https://downloads.apache.org/spark/spark-3.5.0/spark-3.5.0-bin-hadoop3.tgz tar -xzf spark-3.5.0-bin-hadoop3.tgz export SPARK_HOME=$(pwd)/spark-3.5.0-bin-hadoop3 export PATH=$SPARK_HOME/bin:$PATH提示:
hadoop3后缀表示该包内置Hadoop 3.x client,能直接读写HDFS、S3、OSS等,比hadoop2.7版本兼容性更好。如果你后续要连MinIO或阿里云OSS,这个后缀是刚需。
验证Spark是否识别Java:
$SPARK_HOME/bin/spark-shell --version # 正确输出应包含:Spark version 3.5.0, Using Scala version 2.12.18 # 如果报错"Unable to load native-hadoop library",忽略——本地模式不依赖Hadoop native库2.2 用spark-shell交互式跑通WordCount:观察Driver日志里的Stage拆分逻辑
别急着写Python脚本。先用Scala在spark-shell里手敲一遍,因为它的REPL会实时打印DAG可视化和Stage执行日志,这是理解“为什么分Stage”的唯一捷径:
// 启动shell,指定本地模式(单线程)便于观察 $SPARK_HOME/bin/spark-shell --master local[1] // 粘贴以下代码(注意:路径替换成你本地的真实txt文件) val textFile = spark.read.textFile("/home/user/data/sample.txt") val counts = textFile .flatMap(line => line.split(" ")) // Stage 0: mapPartitions at <console>:24 .filter(word => word.nonEmpty) // 同Stage 0:窄依赖,可合并 .map(word => (word.toLowerCase, 1)) // 同Stage 0 .reduceByKey(_ + _) // Stage 1: shuffleMapTask!关键分界点 counts.show(10)观察控制台输出:
Stage 0日志末尾有Number of tasks = 1(因为local[1]只启1个线程)Stage 1开始前会打印Shuffle is enabled,并生成shuffle_0_0_0.index临时文件counts.show()触发Action,此时才真正执行计算
参数说明:
local[1]表示仅用1个CPU核心模拟单节点;换成local[*]会用满所有核,但日志会被冲刷,不利于初学观察。reduceByKey强制Shuffle,而groupByKey会更慢——这是初级实践必须建立的第一个性能直觉。
2.3 用spark-submit提交打包脚本:验证脱离REPL的完整生命周期
交互式只是起点。生产中所有任务都走spark-submit,它会启动独立JVM进程,隔离Driver与Executor内存。写一个最小化可提交脚本wordcount.py:
# wordcount.py from pyspark.sql import SparkSession from sys import argv spark = SparkSession.builder \ .appName("WordCountLocal") \ .master("local[2]") \ # 显式声明2核,避免默认local[*]引发资源争抢 .getOrCreate() lines = spark.read.text(argv[1]) # 从命令行读取输入路径 words = lines.rdd.flatMap(lambda line: line[0].split()) \ .filter(lambda w: len(w) > 0) \ .map(lambda w: (w.lower(), 1)) \ .reduceByKey(lambda a, b: a + b) words.coalesce(1).write.mode("overwrite").csv(argv[2]) # 强制合并为1个输出文件,方便查看 spark.stop()提交命令(注意路径必须绝对):
$SPARK_HOME/bin/spark-submit \ --master local[2] \ --driver-memory 2g \ wordcount.py /home/user/data/sample.txt /home/user/output/wordcount验证输出:
head -n 5 /home/user/output/wordcount/part-00000-*.csv # 应看到类似:(apple,12), (banana,8), ...关键细节:
--driver-memory 2g不是可选——当处理10MB以上文本时,Driver需加载整个DAG并协调Shuffle,内存不足会导致OutOfMemoryError: GC overhead limit exceeded。本地模式下,Driver内存=整个JVM堆内存,必须显式设。
3. RDD vs DataFrame:什么时候该用rdd.map(),什么时候必须切到df.filter().groupBy()?
Spark初级实践最大的认知陷阱,是把RDD和DataFrame当成“两种写法”。实际上,它们是不同抽象层级的执行载体:RDD暴露分区、序列化、Shuffle细节,适合做ETL脏数据清洗;DataFrame隐藏物理计划,靠Catalyst优化器自动剪枝、谓词下推、代码生成,适合结构化分析。选错API,轻则性能差3倍,重则写出无法在集群运行的代码。
3.1 用RDD处理非结构化日志:解析带嵌套空格的Nginx访问日志
假设你拿到的原始日志长这样(access.log):
192.168.1.100 - - [10/Jan/2024:12:34:56 +0800] "GET /api/v1/users?id=123 HTTP/1.1" 200 1234 "-" "Mozilla/5.0"字段间空格数不固定,正则解析易错。RDD的灵活性在此刻体现:
from pyspark import SparkContext sc = SparkContext("local[2]", "NginxParser") def parse_nginx_line(line): try: parts = line.split(' ', 5) # 只切前5段,第6段保留完整request ip = parts[0] timestamp = parts[3][1:] # 去掉开头[ request = parts[5].split(' ')[0] if len(parts) > 5 else "" return (ip, timestamp, request) except: return ("", "", "") # 脏数据兜底 rdd = sc.textFile("/home/user/data/access.log") \ .map(parse_nginx_line) \ .filter(lambda x: x[0] != "") \ .map(lambda x: f"{x[0]}\t{x[1]}\t{x[2]}") rdd.saveAsTextFile("/home/user/output/parsed_access")为什么不用DataFrame?因为
spark.read.text()会把整行当一列,后续df.withColumn("ip", col("value").substr(1,15))需要硬编码位置,一旦日志格式微调就全崩。RDD让你用Python原生逻辑逐行处理,可控性强。
3.2 用DataFrame分析结构化农产品价格CSV:发挥Catalyst优化器威力
对比场景:你有vegetable_prices.csv,含列date,city,commodity,price,unit,需统计“北京每日蔬菜均价”。此时DataFrame是唯一合理选择:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, avg, to_date spark = SparkSession.builder \ .appName("VegPriceAnalysis") \ .master("local[2]") \ .getOrCreate() # 自动推断schema(生产环境务必显式定义schema,避免类型错误) df = spark.read.option("header", "true") \ .option("inferSchema", "true") \ .csv("/home/user/data/vegetable_prices.csv") # Catalyst会自动将filter下推到读取阶段,跳过不满足条件的文件分片 result = df.filter(col("city") == "北京") \ .withColumn("date", to_date(col("date"), "yyyy-MM-dd")) \ .groupBy("date") \ .agg(avg("price").alias("avg_price")) \ .orderBy("date") result.show(5) # 输出:+----------+------------------+ # | date| avg_price| # +----------+------------------+ # |2024-01-01| 5.230000000000001| # |2024-01-02|5.1899999999999995| # +----------+------------------+关键洞察:
.filter(col("city") == "北京")这行代码,在物理执行计划里会变成FileScan csv [date#1, city#2, commodity#3, price#4, unit#5],其中PushedFilters: [*IsNotNull(city), *EqualTo(city,北京)]——意味着Spark SQL引擎在读取CSV时,就跳过了所有city列不为“北京”的行,而不是读全再过滤。这是DataFrame性能碾压RDD的根本原因。
3.3 混合使用RDD与DataFrame:用RDD清洗、DataFrame分析的黄金组合
真实项目永远不是非此即彼。典型工作流:用RDD做脏数据清洗 → 转成DataFrame → 用SQL做聚合。例如农产品数据中混有price="暂无"字符串:
# Step1: RDD清洗,把"暂无"转为None clean_rdd = sc.textFile("/home/user/data/vegetable_prices.csv") \ .map(lambda line: line.split(",")) \ .filter(lambda row: len(row) >= 5) \ .map(lambda row: [ row[0], row[1], row[2], None if row[3].strip() == "暂无" else float(row[3]), row[4] ]) # Step2: 转DataFrame(必须提供schema,否则float列会变string) from pyspark.sql.types import StructType, StructField, StringType, FloatType, DateType schema = StructType([ StructField("date", StringType(), True), StructField("city", StringType(), True), StructField("commodity", StringType(), True), StructField("price", FloatType(), True), # 关键:显式声明为Float StructField("unit", StringType(), True) ]) df = spark.createDataFrame(clean_rdd, schema) # Step3: DataFrame分析(此时price已是数值,可直接avg) df.filter(df.price.isNotNull()).groupBy("city").avg("price").show()血泪经验:
createDataFrame(rdd, schema)比rdd.toDF(schema)更稳定,后者在Spark 3.5.0中偶发类型推断失败。且FloatType()必须显式写,不能用"float"字符串——这是新手常踩的类型陷阱。
4. Shuffle调优实战:为什么你的join慢如蜗牛?3个必调参数让农产品价格关联提速5倍
当你对两个百万级数据集做df1.join(df2, "commodity"),却发现Stage卡在Shuffle Read长达10分钟,CPU利用率却只有15%——这不是代码问题,是Shuffle机制被默认参数拖垮了。Spark的Shuffle不是简单“把数据发过去”,而是涉及序列化、网络传输、磁盘落盘、内存缓存、合并排序五层协作。初级实践必须亲手调这3个参数,否则永远在猜“为什么慢”。
4.1spark.sql.adaptive.enabled=true:让Spark自己决定要不要Shuffle
Spark 3.2+引入自适应查询执行(AQE),能在运行时动态优化Shuffle。对农产品价格关联这种“左表城市维度小、右表价格明细大”的场景,AQE可自动将BroadcastHashJoin替换SortMergeJoin:
spark = SparkSession.builder \ .appName("AQEJoin") \ .config("spark.sql.adaptive.enabled", "true") \ .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \ .config("spark.sql.adaptive.skewJoin.enabled", "true") \ .master("local[2]") \ .getOrCreate() # 城市维度表(<1000行) cities_df = spark.read.csv("/home/user/data/cities.csv", header=True) # 价格明细表(>100万行) prices_df = spark.read.csv("/home/user/data/prices.csv", header=True) # AQE会在运行时判断:cities_df太小,自动广播!无需写broadcast(cities_df) result = prices_df.join(cities_df, "city_id") \ .filter("price > 10") \ .groupBy("province").count() result.show()验证AQE生效:看Spark UI的SQL tab,Execution Plan里会出现
AdaptiveSparkPlan isFinalPlan=true,且BroadcastHashJoin节点旁标注Broadcasted。若没出现,检查cities_df.count()是否真小于spark.sql.autoBroadcastJoinThreshold(默认10MB,可调大)。
4.2spark.sql.adaptive.localShuffleReader.enabled=true:用本地磁盘加速Shuffle读取
默认Shuffle读取走网络(即使本地模式),而localShuffleReader让Executor优先从本地磁盘读Shuffle文件,减少网络抖动:
# 在spark-submit中添加 --conf spark.sql.adaptive.localShuffleReader.enabled=true \ --conf spark.sql.adaptive.coalescePartitions.enabled=true \实测效果(本地模式,2GB价格数据):
| 配置 | Shuffle Read时间 | CPU利用率 |
|---|---|---|
| 默认 | 42s | 35% |
| 启用localShuffleReader | 18s | 82% |
原理:
localShuffleReader绕过Netty网络栈,直接mmap本地Shuffle文件。但注意——它只在local[*]或Standalone集群有效,YARN/K8s需额外配置spark.shuffle.service.enabled=false。
4.3spark.sql.files.maxPartitionBytes=128m:控制Shuffle分区大小,避免小文件风暴
当农产品价格数据按日期分区存储(/data/prices/year=2024/month=01/day=01/),Spark默认按128m切分每个文件。但若某天数据只有2MB,就会生成64个超小分区,Shuffle时产生海量小文件,拖垮磁盘IO:
# 读取时强制合并小文件 df = spark.read.option("maxPartitionBytes", "128m") \ .option("recursiveFileLookup", "true") \ .csv("/home/user/data/prices/") # 或提交时全局设置 --conf spark.sql.files.maxPartitionBytes=128m \参数说明:
maxPartitionBytes不是“目标大小”,而是“上限”。Spark会尽量让每个分区≤该值,但不会把130MB文件硬切成两半(会保持整文件)。对小文件多的场景,设为256m或512m更合理。
5. 避坑指南:Spark初级实践最常踩的5个坑,现象、原因、解决全写清楚
新手在Spark初级实践里摔的跟头,90%集中在环境、API、Shuffle三类。下面5条是我带过的27个实习生、3个外包团队反复验证过的“血坑”,每条都按“现象→原因→解决”写透,不讲虚的。
5.1 现象:spark-shell启动报错java.lang.NoClassDefFoundError: scala/Product
原因:Spark二进制包自带Scala 2.12,但你的系统SCALA_HOME指向了2.13或2.11,导致类加载冲突。
解决:彻底删除SCALA_HOME环境变量,Spark会用自己的Scala。验证:$SPARK_HOME/bin/spark-shell --version输出中Using Scala version 2.12.18必须与包名一致(spark-3.5.0-bin-hadoop3.tgz对应2.12)。
5.2 现象:df.write.csv()报错org.apache.hadoop.security.AccessControlException: Permission denied: user=xxx, access=WRITE, inode="/user"
原因:本地模式下Spark仍会尝试连接HDFS,因未配置core-site.xml,默认连到hdfs://localhost:9000,而你根本没起HDFS。
解决:强制指定本地文件系统——在spark-submit加参数:--conf spark.hadoop.fs.defaultFS=file:///。或者写CSV时用绝对路径:df.write.csv("file:///home/user/output/")。
5.3 现象:df.join()后count()返回0,但单独查左右表都有数据
原因:Join字段类型不一致。例如左表city_id是StringType,右表是IntegerType,Spark静默转为null导致匹配失败。
解决:用df.printSchema()检查两边字段类型,强制统一:df1.withColumn("city_id", col("city_id").cast("string"))。永远不要信inferSchema。
5.4 现象:本地跑spark-submit正常,提交到YARN集群报Container exited with a non-zero exit code 143
原因:YARN的yarn.nodemanager.vmem-pmem-ratio默认2.1,即虚拟内存不得超过物理内存2.1倍。Spark Executor的-Xmx设了4g,但JVM额外开销使虚拟内存达10g,被YARN Kill。
解决:提交时加--conf spark.yarn.executor.memoryOverhead=2048(单位MB),或调高YARN参数(需集群权限)。
5.5 现象:df.filter("price > 10").count()比df.filter(col('price') > 10).count()慢3倍
原因:SQL字符串过滤(filter("price > 10"))无法触发Catalyst谓词下推,Spark先读全量数据再过滤;而col()方式能生成优化后的物理计划。
解决:永远用col("field")或df["field"],禁用字符串SQL过滤。这是Spark 3.x的硬性最佳实践。
6. 进阶技巧:用spark.ui.port和spark.eventLog.dir把本地调试变成可回溯的黑匣子
初级实践最痛苦的,不是写不出代码,而是报错时不知道哪一行触发了哪个Stage、Shuffle写了多少文件、GC停顿了几次。Spark UI和事件日志就是你的黑匣子记录仪。但默认配置下,它们要么打不开,要么日志删得比你反应还快。下面这套配置,让我能把每次spark-submit的完整执行过程存档,回溯任意一次失败。
6.1 永久开启Spark UI并绑定固定端口,避免端口冲突
Spark UI默认随机端口(如4040、4041),当你同时跑多个spark-shell,UI会抢占失败。在$SPARK_HOME/conf/spark-defaults.conf里加:
spark.ui.port 4040 spark.ui.retainedStages 100 spark.ui.retainedJobs 100 spark.ui.retainedApplications 100然后启动时指定历史服务(让UI不随Driver退出而消失):
# 启动历史服务(后台运行) $SPARK_HOME/sbin/start-history-server.sh # 提交任务时启用事件日志 $SPARK_HOME/bin/spark-submit \ --conf spark.eventLog.enabled=true \ --conf spark.eventLog.dir=file:///home/user/spark-events \ --master local[2] \ wordcount.py /input /output验证:浏览器打开
http://localhost:4040,能看到当前运行的所有Application。即使任务结束,历史服务也会从/home/user/spark-events加载日志,显示完整的DAG、Stage耗时、Shuffle读写量。
6.2 解析事件日志定位Shuffle瓶颈:用spark-sql查日志本身
Spark事件日志是JSON格式,但直接用cat看是灾难。Spark自带spark-sql可直接查询日志:
# 启动spark-sql并加载事件日志 $SPARK_HOME/bin/spark-sql \ --conf spark.sql.adaptive.enabled=false \ --conf spark.sql.adaptive.coalescePartitions.enabled=false \ -e "CREATE TEMPORARY VIEW event_log USING json OPTIONS (path '/home/user/spark-events');" # 查看Shuffle写入最多的Stage spark-sql> SELECT event, `Stage Info`.`Stage ID`, `Stage Info`.`Number of Tasks`, `Stage Info`.`RDD Info`.`Storage Level` FROM event_log WHERE event = 'SparkListenerStageCompleted' ORDER BY `Stage Info`.`RDD Info`.`Memory Bytes Spilled` DESC LIMIT 5;输出示例:
Stage ID: 3,Number of Tasks: 200,Memory Bytes Spilled: 1247890123—— 这说明Stage 3的200个Task共溢出1.2GB到磁盘,是性能瓶颈。此时你应该检查该Stage的reduceByKey或join操作,调大spark.sql.adaptive.coalescePartitions.enabled或增加分区数。
6.3 用spark.metrics.conf导出JVM指标到本地文件,监控GC压力
Driver和Executor的GC停顿是隐形杀手。在$SPARK_HOME/conf/metrics.properties中配置:
*.sink.file.class=org.apache.spark.metrics.sink.FileSink *.sink.file.period=10 *.sink.file.unit=seconds *.sink.file.directory=/home/user/spark-metrics启动后,/home/user/spark-metrics下会生成metrics-*.json,内容含:
{ "jvm.pools.Metaspace.usage.used": 124567890, "jvm.gc.PS-MarkSweep.time": 12345, "jvm.gc.PS-Scavenge.time": 6789 }实操价值:当
PS-MarkSweep.time持续>10000ms(10秒),说明老年代GC频繁,需调大--driver-java-options "-XX:MaxMetaspaceSize=512m";若PS-Scavenge.time突增,是年轻代不够,加-Xmn2g。
我坚持给每个Spark任务配UI和事件日志,不是为了炫技,而是因为在集群上,你永远不知道下一个OOM发生在哪个Executor的第几秒。把本地调试变成可回溯的黑匣子,是初级实践通往可靠交付的最后一步。希望帮到你。
本文还有配套的精品资源,点击获取