Spark实验七实战:RDD编程与Spark SQL读JSON入门
简介这份实验报告围绕Spark初级编程实践展开面向正在学习大数据、Hadoop与Spark课程的高校学生完整呈现从环境搭建到独立应用开发的全过程。报告基于Windows 10宿主机与Ubuntu Kylin 16.04虚拟机环境给出Hadoop 3.1.3、JDK 1.8的具体安装配置方法并详细记录了使用./bin/spark-shell启动Spark、通过sc.textFile读取本地文件与HDFS文件并统计行数、编写Scala独立应用并利用sbt打包、通过spark-submit提交到集群运行的完整操作链路。除基础流程外还包含对学生成绩求平均值的数据处理示例附有SimpleApp、RemDup、AvgScore三个应用的完整代码以及实验过程中常见异常如IllegalArgumentException、InvalidInputException、URISyntaxException的成因分析与解决办法。资源为单个docx文档共1份压缩包大小1.9MB内容结构清晰既适合实验前预习也可作为实验报告撰写与排错参考。该资源已有8338人学习下载性价比较高。1. 拿到实验七先别急着背 RDD 算子第一次做《实验七Spark初级编程实践》的人最常见的翻车现场不是算子写错而是环境装了、包也导了一个reduceByKey却卡在原地转圈或者程序跑完结果和预期完全对不上。这个实验要解决的正是从“环境就绪”到“程序跑通”之间的断层核心就三件事选对运行模式、写通一个最小 RDD 程序、再尝试用 Spark SQL 把半结构化数据拉进来。面向的是刚接触分布式计算、想用最短路径验证 Spark 能干什么的初学者。这篇笔记会把部署参数、代码细节、故障排查逐一拆开让你照着走一遍就能把作业交得干净也顺便理解这个“分布式计算框架”到底是怎么工作的。2. Spark 环境选型为什么 local 模式足够做完实验七及四个必调参数2.1 三种部署模式对比local、standalone 与 yarn 的分工很多初学者一上来就纠结要不要搭集群甚至去查“Spark集群搭建”的教程结果折腾三天网卡、内存、权限全出问题作业还没开始写。实验七题目里带“初级”两个字就意味着它的教学目标不是让你部署生产环境而是让你理解 RDD 计算模型。我用过三种模式给你的建议很直接本地开发用local作业验收用local只有老师明确要求“搭建集群”时才去碰 standalone。三者的本质区别在于资源调度和进程分布。local模式下Spark 的所有组件——driver、executor、调度器——都跑在同一个 JVM 进程里只是用多线程模拟并行计算standalone 模式是 Spark 自带的集群管理器会有独立的 master 节点和 worker 节点你需要配spark-env.sh、起进程、盯端口yarn 模式则是把资源统一交给 Hadoop 的 ResourceManager 调度适合公司里已经有 Hadoop 集群的场景。对实验七来说local模式跑出来的 RDD 计算逻辑和集群上完全一致区别只在并行度和数据量上限。对比项localstandaloneyarn部署难度解压即用需配置 master/worker需先有 Hadoop 集群适用场景开发、测试、课程实验小规模自建集群生产环境、与 HDFS 配合进程分布单 JVM 多线程多节点多进程由 YARN 统一调度调试便利性高日志集中在控制台中需要查多个节点日志低需要看容器日志我一般会这样建议如果你是 Windows 笔记本 Jupyter直接local如果你在 Linux 虚拟机里想感受一下集群的 Web UI再搭 standalone。不要在课程实验阶段同时挑战“集群搭建”和“Spark 编程”两件事那样出了故障都分不清是环境问题还是代码问题。2.2 spark-submit 提交脚本时的四个必调参数环境选型确定后下一步是提交方式。你当然可以在 pyspark 交互式 shell 里一行行敲代码但实验报告通常要求“脚本 运行结果”所以我更推荐用spark-submit提交一个.py文件。下面是实验七最常用的最小提交命令spark-submit \ --master local[2] \ --driver-memory 2g \ --executor-memory 2g \ --conf spark.sql.shuffle.partitions4 \ wordcount.py--master local[2]里的方括号数字是最容易被忽略的参数。它代表本地模拟的 Executor 线程数直接决定任务的并行度。写local[*]会使用机器全部 CPU 核心但如果你同时在跑浏览器、IDE 和虚拟机机器会明显卡顿。实验阶段设 2 或 4 就够数据量小线程多了反而在调度上浪费时间。--driver-memory和--executor-memory在 local 模式下实际作用于同一个 JVM但分着写能让你理解 driver 和 executor 是两类角色——前者负责任务规划和结果汇总后者负责真正执行计算。机器内存是 8GB 的话两个都设 2g 比较安全16GB 的机器可以都设 4g。我曾见过有人把--executor-memory设成 16g然后系统直接 OOM这不是 Spark 的问题是对“内存”没有边界意识。--conf spark.sql.shuffle.partitions4很多人不写因为默认值也有 200。问题在于默认的 200 个分区是针对大数据的实验数据只有几 MB200 个分区纯粹是额外开销。设成 4配合local[2]既能保证每个分区被处理又不会产生上百个空任务。一句话总结四个参数里master 决定并行方式两个 memory 决定内存上限shuffle.partitions 决定数据打散粒度。2.3 用 pyspark 交互式 shell 做冒烟测试环境通不通两行命令验证正式写脚本前我强烈建议先做个三十秒的冒烟测试确认 Spark 本身能正常工作。打开终端输入pyspark进入交互式环境然后执行lines sc.textFile(file:///tmp/test.txt) lines.count()如果test.txt里随便写了三行文字这里会返回3说明环境通。两个细节需要留意一是读本地文件必须加file:///前缀不加的话 Spark 默认去 HDFS 找路径会报FileNotFoundException二是sc是 pyspark 交互环境里已存在的 SparkContext 对象不用自己创建这也是初学者最容易懵的地方。冒烟测试通过后打开浏览器访问http://localhost:4040你会看到 Spark Web UI。这个页面是后续所有性能判断的“黑匣子”窗口——Jobs、Stages、Executors 三个标签页分别记录任务进度、执行阶段和资源占用。实验阶段不用看懂全部只需要确认运行count()时页面上出现了一个 Completed Job就说明提交链路没问题。这个动作花不到一分钟却能帮你把“环境问题”和“代码问题”切开。3. RDD 初级编程从并行集合到 WordCount 的完整落地3.1 创建 RDD 的三种途径parallelize、textFile 与 wholeTextFiles实验七的核心是 RDD 编程而所有 RDD 程序的第一步都是“把数据变成 RDD”。常见做法有三种parallelize适合把内存中的集合转成分布式数据集比如测试时临时生成一批数字textFile用于读取文本文件按行切分返回的每个元素就是一行字符串wholeTextFiles按文件切分返回的是(文件名, 文件内容)的键值对。# 方式一内存集合并行化 nums sc.parallelize([1, 2, 3, 4, 5], 2) # 方式二读取文本文件注意 file:/// 前缀 lines sc.textFile(file:///data/words.txt) # 方式三整文件读取适合处理一堆小文件 file_rdd sc.wholeTextFiles(file:///data/txt_dir/)parallelize的第二个参数2是分区数表示数据会被切成 2 份。分区数直接决定并行度上限——分区数小于 Executor 核心数时部分核心会闲置大于核心数时会有额外的任务调度开销。实验数据下分区数设 2 或 4 足够不用刻意追求多。需要提醒的是wholeTextFiles会把整个文件内容当作一个字符串载入内存处理大文件时极易 OOM初级实验里只适合处理几 KB 的小文件目录。创建完 RDD 后可以用getNumPartitions()验证分区情况这是排查并行度问题最直接的手段。如果返回的份数和预期不符检查spark.default.parallelism配置是否被显式覆盖或者 master 模式是否写成了单线程的local[1]。3.2 Transformation 与 Action 的本质区别为什么 map 之后终端没任何输出RDD 编程最容易让新手困惑的问题执行了map、filter总算在想“程序是不是卡住了”。这不是卡住而是 Spark 的懒惰执行特性。RDD 算子分两类Transformation 只定义计算步骤不真正计算Action 才会触发实际执行并把结果返回 driver 或写到外部存储。lines sc.textFile(file:///data/words.txt) upper_lines lines.map(lambda x: x.upper()) # 这里没有任何执行发生 count upper_lines.count() # 执行到这里Spark 才开始真正计算map属于 Transformation它只是建立了一个“血缘关系图”告诉 Spark 从lines开始经过一个匿名函数得到新的数据集。真正开始执行的是count()——这是一个 Action。理解这一点极其重要因为实验报告里考察的往往不是你会不会调 API而是你能否解释“为什么 map 之后没有输出”。另一个需要区分的点是reduceByKey与groupByKey。前者会在 map 端先做一次局部聚合然后只把聚合后的中间结果 shuffle 到下游后者会把同一个 key 的全部分组数据都传输到下游。相同逻辑下reduceByKey的网络开销小得多这也是 Spark 性能调优的第一课。实验七里做 WordCount 时务必用reduceByKey这也是 Spark 社区和面试题里反复强调的差别。3.3 WordCount 完整脚本与结果收集的取舍下面给你一个可直接提交运行的 WordCount 完整脚本注释写清楚每个步骤对应 Spark 的哪个阶段from pyspark import SparkContext, SparkConf # 初始化setAppName给应用起名setMaster指定本地运行模式 conf SparkConf().setAppName(E7_WordCount).setMaster(local[2]) sc SparkContext(confconf) # 读取输入文件按行切分 lines sc.textFile(file:///data/words.txt) # 切分单词flatMap将每行的单词列表“压平”filter去掉空白词 words lines.flatMap(lambda line: line.split( )).filter(lambda w: len(w) 0) # 每个单词映射为 (word, 1)再按 key 做加法聚合 pairs words.map(lambda w: (w, 1)) counts pairs.reduceByKey(lambda a, b: a b) # 按词频倒序排列并取出前10个结果 top_n counts.sortBy(lambda x: x[1], ascendingFalse).take(10) # 打印结果 for word, count in top_n: print(f{word}: {count}) sc.stop()这段代码里有几个设计值得说明。第一flatMap而不是map是因为split( )的返回值是一个列表map会得到“列表的列表”flatMap会将二级列表压平得到单词流。第二filter去掉空字符串是必要的连续空格会导致空串参与计数。第三结果用take(10)而不是collect()因为collect会把 RDD 的全部分区数据拉回 driver实验作业如果输入文件是几百 MBcollect会让 driver 内存直接爆掉。take只取前 N 条安全且满足查看效果的需求。跑完这个脚本你会得到按词频降序的单词列表。如果结果为空不要查算子先看words.txt的路径前缀和文件编码如果结果比预期多基本是标点符号没做清洗比如hello,和hello被当成两个单词。这是 WordCount 实验最经典的结果偏差来源处理方式是在切分后用正则去掉非字母字符。4. 用 Spark SQL 读取 JSON把初级实验从 RDD 做厚到结构化查询4.1 SparkSession 与 SparkContext为什么读 JSON 前要换入口实验七如果只做 WordCount其实只覆盖了 Spark 的 RDD 接口。但热词里“spark sql”和“spark中读取json”出现频率很高说明很大一部分老师会把“读取半结构化数据”也塞进实验要求里。读取 JSON 这类有 schema 的数据正确工具不是 SparkContext而是 SparkSession。SparkSession 是 Spark 2.0 之后统一的入口它内部封装了 SparkContext、SQLContext 和 HiveContext。简单理解SparkContext 面向无结构的 RDDSparkSession 面向结构化数据。你用 recall 写spark.read.json(...)时本质上是在调用 SQL 引擎的能力。实验脚本里这样初始化from pyspark.sql import SparkSession # 创建 SparkSession同样支持本地模式 spark SparkSession.builder \ .appName(E7_JSONDemo) \ .master(local[2]) \ .getOrCreate()注意这里不再手动创建 SparkContext。getOrCreate()的含义是如果当前 JVM 里已有 SparkContext直接复用如果没有再创建新的。这在同一个环境里交叉使用 RDD 和 DataFrame 时很重要——重复创建会导致SparkContext already exists的报错。4.2 read.json 的 Schema 推断与 multiLine 参数Spark 读取 JSON 时有自动推断 schema 的能力它会扫描数据推断出每列类型。大多数场景下这是很方便的但它也有代价对小文件需要额外启动一次采样扫描。对于实验数据直接使用默认推断即可。以下代码演示spark.read.json读取本地一个 JSON 文件# 读取单行 JSON 文件 df spark.read.json(file:///data/students.json) # 打印 schema 信息可用于检查类型推断是否正确 df.printSchema() # 查看前 5 行数据 df.show(5)这里最常踩的坑是“JSON 文件是多行格式”。默认情况下 Spark 把“文件的一行”当作一个 JSON 对象如果你的文件里每个对象拆成多行比如带格式化缩进的 JSON 数组就会报JSON parsing failed。解决办法是加multiLine参数df spark.read.option(multiLine, True).json(file:///data/students_multiline.json)参数名默认值作用何时修改multiLinefalse控制是否将整个文件当作一条 JSON文件有格式化缩进时设为 trueinferSchematrue自动推断列类型数据量大且 schema 已知时可设为 falseprimitivesAsStringfalse将基本类型全部读成字符串后续处理需要统一字符串时使用modePERMISSIVE解析失败时容错策略实验建议用FAILFAST快速暴露错误printSchema()是读 JSON 后必做的动作。它会列出每列的名称、类型和是否可空。实验里拿到的 JSON 经常出现“age 字段有时是整型、有时是字符串”的情况不打印 schema 很难发现。如果发现类型推断不符合预期可以在read.json前用.schema(name STRING, age INT)显式指定避免 Spark 猜错。4.3 注册临时视图后用 SQL 查询把 SQL 基础直接迁移过来读取到 DataFrame 之后实验常见的后续操作包括选择列、过滤、分组、聚合。DataFrame 支持方法式操作也支持注册为临时视图后用 SQL 语句。我一般更推荐后者因为对初学者来说SQL 是可迁移的已有技能写起来更顺手且实验报告展示时也更直观。# 读取 JSON 并注册为临时视图 df.createOrReplaceTempView(students) # 用纯 SQL 做分组聚合查询 result spark.sql( SELECT department, COUNT(*) AS cnt, AVG(score) AS avg_score FROM students WHERE score 60 GROUP BY department ORDER BY avg_score DESC ) # 触发计算并打印结果 result.show()createOrReplaceTempView创建的是一个仅当前 SparkSession 可见的临时表它不会留在磁盘进程结束就消失。查询里比较重要的是WHERE的执行时机——它会在聚合前过滤数据这能显著减少 shuffle 的数据量。Spark SQL 的查询优化器会做谓词下推但实验阶段你不需要理解优化器细节只需知道先过滤、再聚合、最后排序这个书写习惯符合 Spark 的期望。这段代码和 RDD WordCount 形成对比RDD 强调的是“怎么把计算步骤拆成分区上的操作”Spark SQL 强调的是“描述你要什么结果引擎负责决定怎么算”。实验七能把这个对比讲清楚报告就已经超出大多数同类作业的水准。5. 常见问题排查让新手卡在原地的 5 个 Spark 故障5.1 8080 端口打不开 Web UI现象照着教程访问http://localhost:8080页面一直打不开但 Spark 程序明明在运行。原因8080 是 standalone 模式下 Master 的 Web UI 端口你用的是 local 模式根本没有独立的 Master 进程自然没有 8080 服务。解决local 模式下要看 Spark UI访问http://localhost:4040。这是 driver 的端口只要有一个 SparkContext 在运行这个页面就会存在。如果你开了多个 SparkSession 或重复运行了 pyspark端口会顺延为 4041、4042。想看多个任务界面注意区分 URL 里的端口号。5.2 程序跑了一半突然报错退出日志里全是 ClassNotFound现象本地跑无异常换了一台机器或者交给老师检查时运行到某个算子突然报java.lang.NoClassDefFoundError重启无用。原因最常见的是 JDK 版本问题。Spark 3.x 对 Java 8 支持最完善如果你用的默认 JDK 是 17 或 21部分组件尤其是涉及动态代码生成的算子会触发兼容问题表现就是运行到一半才崩溃。解决先执行java -version确认版本。如果是 17切换到 JDK 8 再试。Windows 下安装多个 JDK 时设置JAVA_HOME和PATH时容易遗漏控制台会话的更新——改完环境变量要重开终端否则 Spark 读取到的还是旧版本。这个问题很玄学经常让人怀疑是代码问题实际上就是版本错配。5.3 collect() 把整份数据拉回 driver直接内存溢出现象WordCount 或 JSON 查询在.show()之前一切正常加上.collect()后立刻 OOM屏幕输出一段堆栈后进程退出。原因collect()是 Action它要求所有 Executor 把计算结果全量传输到 driver 进程。实验机器内存只有几个 Gdriver 自然扛不住。解决把collect()换成take(n)或show(n)只取前 N 行。如果非要全量数据做后续分析就用df.write.saveAsTable(...)或coalesce(1).write.csv(...)把结果写入文件再单独读文件。记住一句话collect 是危险动作数据量大时它是第一杀手。5.4 Executor 启动失败报 Python 相关错误现象本地 local 模式没问题但只要用spark-submit --master local[*]提交 Python 文件Executor 就报Python worker failed to connect back。原因PySpark 的 Executor 需要启动一个 Python 进程它通过环境变量PYSPARK_PYTHON找到 Python 解释器。如果你电脑里装了多个 Python比如 Anaconda 和系统自带的Executor 可能指向了一个不存在的路径或错误版本。解决提交前显式指定export PYSPARK_PYTHONpython3如果你用虚拟环境就填虚拟环境的完整路径。注意在spark-submit前执行这个 export 命令且在同一终端会话内生效。后续每次提交都要确认这个变量没被重置。5.5 读中文 JSON/CSV 输出乱码现象用spark.read.json读取中文数据show()结果里的中文字段全是问号或乱码英文正常。原因文件编码和 Spark 默认读取编码不一致。Spark 默认按 UTF-8 读取如果你的文件是 GBK 编码或者爬虫生成的数据混入了非 UTF-8 字符就会乱码。另一种情况是终端本身编码不对数据没问题但显示乱。解决先用file命令或文本编辑器确认文件编码如果是 UTF-8 带 BOM用option(encoding, utf-8)读取如果源头是 GBK先转码iconv -f GBK -t UTF-8 file.csv file_utf8.csv再交给 Spark。实验作业里提交的数据集最好统一转成 UTF-8这也是一个能写进实验报告的数据预处理步骤。6. 用一个“预期对照表”验证结果并把 Spark UI 当体检报告6.1 手写预期输出让程序结果有对照实验做完了怎么确认结果是“对”的很多人只看一眼结果不为空就觉得成功了这不够。我养成的习惯是先构造一份极小的样本数据手算预期结果再跑程序对照。比如准备一个只有 5 行、20 个单词的small_words.txt手动数出每个单词出现次数做成对照表单词手算次数程序输出是否一致spark33是hadoop22是sql55是这样做的价值在于它能同时验证三件事数据读取路径是否正确、切分逻辑是否符合预期、聚合逻辑有没有写错。WordCount 这类程序如果直接上真实数据集你根本不知道输出错的还是在错的。先小样本验证再换大数据跑正式结果这条路径我至今沿用。6.2 从 Spark UI 的 Executors 与 SQL 页面确认资源真的用上了程序跑通之后还有最后一步值得做打开localhost:4040的 Executors 页面确认任务确实分布式执行了而不是所有数据在一个线程里串行处理。观察两个指标Executor 的数量是否和你--master local[2]设置的并行度一致以及“Shuffle Read”和“Shuffle Write”的数据量是否在合理范围。如果你看到 Shuffle Read 几 GB 而输入文件只有几 MB说明分区设置或 key 分布存在问题如果 Shuffle 数据量为 0可能reduceByKey根本没触发 shuffle结果大概率不对。SQL 页面则能查到每个查询的执行计划包括过滤条件下推和聚合策略。实验阶段不需要逐行读计划只需确认“Scan json”和“Exchange”两个节点存在前者说明读到了数据后者说明发生了网络数据传输。这两个页面就是 Spark 的体检报告把它和程序结果对照着看比单看控制台日志可靠得多。我做实验七留下的习惯是每跑一个任务都习惯性看一眼 Spark UI确认任务确实由 Executor 执行、shuffle 数据规模与输入规模匹配然后再把结果写进实验报告。这个过程不花费额外功夫却能把“跑通了”的结论落到实处。希望帮到你。本文还有配套的精品资源点击获取