
说实话Spark的知识点多且杂一到复习的时候最容易陷入的状态就是概念都眼熟代码能抄但换个场景就懵一遇到OOM更是一头雾水。我见过不少同学把RDD、DataFrame、Spark Streaming三个词背得很熟但问到底层一个Job是怎么被切成Stage的、谁在负责调度就说不清了。这篇复习汇总不准备面面俱到而是把Spark里最容易被考到、也最影响实战的几块核心知识点重新捋一遍RDD与作业调度、核心算子与实战案例、Spark SQL高频函数与join语义、集群部署与内存调优。不管你是正在准备课程考试还是打算面试大数据岗或者临时要接一个数据清洗的活都可以按这条主线来过一遍。1. 复习前先搭好框架Spark到底在学什么1.1 从MapReduce到Spark它到底改变了什么很多教材会从“Spark是一个分布式的、基于内存的、统一的数据分析引擎”说起但我觉得直接对比MapReduce更容易建立记忆锚点。MapReduce时代每个阶段之间都要把中间结果写到HDFS下一阶段再重新读一遍。如果是那种迭代式的算法模型每一轮训练都要经历一轮落盘和读盘磁盘IO会成为整个作业的瓶颈这是MapReduce最让人难受的地方。Spark的思路就是把“中间结果尽量留在内存里”这件事变成默认动作。它把一个作业描述成有向无环图算子之间有依赖关系调度器能知道哪些计算可以串起来执行哪些必须等待上一个阶段完成。有了这种能力像PageRank、ALS推荐、多轮join这类适合迭代的场景运行时间能够比MapReduce快出一个数量级。当然Spark虽然常被说成“内存计算框架”但它并不是所有场景都适合。如果你只是做一次简单的ETL数据量又不算大用Spark反而要应付集群启动、资源分配这些开销。所以复习的时候先建立“什么场景适合用”的判断力比死记概念更能应对考试里的场景题。1.2 一条主线串起Spark全家桶我复习时最喜欢用一条主线来串联知识点把Spark看成一个“数据流水线”从数据如何读进来开始到数据经过哪些转换再到什么时候真正触发计算结果如何输出中间夹杂着分区、缓存、任务调度、失败重试这些细节。顺着这条线走你自然会问出一系列问题而这些问题的答案恰好就是考试重点。这条主线也对应着Spark生态里的几个核心模块。Spark Core负责最基础的RDD、调度、内存管理是所有功能的地基Spark SQL在Core之上提供DataFrame/Dataset和SQL语法内置Catalyst优化器Spark Streaming基于微批模型做流式处理MLlib里放着ALS、随机森林这些常用算法GraphX则面向图计算。实际项目里用得最多的是Core和SQLStreaming在一些实时场景中出现MLlib的协同过滤模型在电商推荐里很常考GraphX相对冷门一些。1.3 这些知识点到底会在哪些场景被反复问说几个真实常见的项目场景帮你理解复习的方向。网约车订单日志清洗要处理的就是时间字段格式不统一、司机ID和订单ID有重复、部分记录缺少金额这些脏数据全过程基本是Spark SQL的过滤、去重、类型转换、分组聚合。农产品价格数据分析是另一类典型数据源通常是CSV或数据库导出的表格价格字段可能有空值地区名称不统一日期可能既有“20240101”又有“2024-01-01”两种写法。你用Spark做的主要事情就是把数据清洗成统一口径再计算按日的均价、周环比、月同比最后落到BI报表里。电商推荐系统则是面试里的常客离线阶段用ALS算法训练用户-商品隐因子矩阵在线阶段用训练好的模型给用户召回候选商品。这里Spark的MLlib封装了完整的训练流程但面试官更爱问的是Spark分布式计算的底层逻辑比如每个Executor在训练环节里承担了什么、数据在节点间怎么交换。这也是我把RDD和调度机制放在复习第一站的原因后面所有高级API本质上都是在RDD和调度之上加了一层封装。2. RDD与调度机制Spark的骨架不能弯2.1 RDD的本质与五大属性RDD英文全称Resilient Distributed Dataset弹性分布式数据集。考试最喜欢问的第一个问题往往就是“RDD的五大属性是什么”。别只背名词每个属性都对应着一句实现层面的含义。第一分区列表RDD被拆成了多少个分区分区决定了计算的并行度第二计算函数给定一个分区通过什么逻辑算出这个分区的数据第三依赖关系这个RDD是怎么从父RDD算出来的是窄依赖还是宽依赖第四分区器针对key-value型的RDD数据按什么规则分到不同分区第五首选位置在读取文件时每个分区应该优先调度到哪个节点上去执行减少网络传输。如果你觉得五大属性太抽象可以做一个小类比RDD像一批要从仓库发到门店的货分区等价于一个个打包好的箱子计算函数是拆箱规则告诉工人每个箱子里的货怎么归类依赖关系是运输路线图能看出这批货是从哪个仓库调来的中间有没有拆箱重装分区器是配送目的地决定箱子送到哪个门店首选位置则是路线规划告诉司机优先走哪条路最省时间。这个类比不算严谨但能帮你在考场上用最短时间想起每个属性是干什么的。还有一个关键概念叫血缘Lineage。每个RDD都记录了自己是怎么从父RDD演化来的当某个分区的数据丢失后Spark可以根据血缘关系重新计算不需要像MapReduce那样做数据副本的冗余备份这就是“弹性”这个词的由来。复习时一定要区分血缘和Checkpoint血缘是重新计算的依据Checkpoint则是把中间结果真正持久化用来切断太长的血缘链避免节点故障后恢复计算太久。2.2 从Driver到Task一次作业是怎么被切分的考试里另一个高频考点是Application、Job、Stage、Task这四个名词的关系。一个Application对应一个SparkContext你跑一个spark-submit提交的任务就是一个Application。一个Application里可以执行多个Job触发点就是Action算子比如collect、count、saveAsTextFile。一个Job根据宽依赖切成多个StageStage之间的边界通常就是shuffle发生的边界。一个Stage里的数据被分成多个分区每个分区由一个Task去处理Task才是真正在Executor线程池里跑的最小计算单元。实际调度过程是这样的DAGScheduler负责把Job划分成Stage它会把整个依赖图从后往前追踪遇到宽依赖就切一刀切完生成若干Stage然后TaskScheduler把Stage里的Task按照数据本地性原则发布到对应的Executor节点Executor里的TaskRunner接到任务后才真正执行RDD分区上的计算函数。所以你在Spark UI的Jobs页签里看到的每个Stage点进去能看到成千上万个Task每个Task的执行时间、输入数据量、shuffle读写量全都记录在案这也是排查性能问题最直接的窗口。理解宽窄依赖对面试特别重要。窄依赖是指父RDD的每个分区只被子RDD的一个分区使用比如map、filter这种情况下可以做流水线式的合并宽依赖则是父RDD的每个分区可能被子RDD的多个分区使用比如groupByKey、reduceByKey过程必然涉及shuffle也就是把数据按照key重新分区再拉取开销比窄依赖大得多。考试里如果问“为什么Spark要划分Stage”标准答案就是把窄依赖的计算尽可能合并到一起执行只把宽依赖作为Stage的分界点减少shuffle次数。2.3 算子分类、缓存与血缘笔试和面试都爱考Spark算子分两类Transformation是懒加载的它不会立即触发计算只是在构造好这个RDD的谱系关系Action才真正触发计算。这解释了一个经常让新人困惑的现象你写了一大串map、filter、join日志里什么都没有打印Spark UI上也没有Job出现直到你调了count或者saveAsTextFile所有转换才一口气跑起来。这个设计的直接好处是可以做流水线优化比如rdd.map(...).filter(...)会在同一个Stage里被合并执行数据不会在map和filter之间落一次中间结果。常用Transformation里有一组对比要格外注意groupByKey和reduceByKey。两者都能完成按key聚合但groupByKey会先把同一个key的所有value原样收集到一起shuffle的数据量很大reduceByKey会在map端先做一轮预聚合把同一个key的value合并成一个局部结果再通过网络传输数据量急剧减少。实际项目里能用reduceByKey解决的问题尽量不要用groupByKey。缓存和Checkpoint是另外两个必考点。cache或者persist可以避免一个RDD被多个Action反复计算persist可以选择存储级别比如MEMORY_ONLY、MEMORY_AND_DISK、DISK_ONLY。这里有一个常见的坑很多人以为调了cache之后内存就一定能存下全部数据其实如果数据量超过Executor内存cache的数据会被直接丢弃之后用到它时还是要从头算。Checkpoint和cache不同它是真的把RDD数据写到可靠存储里同时切断血缘关系一般建议在shuffle之后做Checkpoint对复杂的迭代作业特别有效。面试时如果被问到“什么时候用cache什么时候用Checkpoint”可以答数据变化不频繁、希望复用中间结果用cache依赖链太长、担心节点故障导致长时间重算用Checkpoint。3. 动手跑通核心案例从WordCount到数据清洗3.1 WordCount的完整实现与隐藏细节WordCount几乎是Spark的“hello world”它麻雀虽小但五脏俱全读数据、做转换、shuffle聚合、输出结果全部都有。我建议不管你是考试还是面试都要能闭着眼写出来并且能解释每一步对应着什么级别的调度动作。Python版本的实现大概是这样的from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(WordCount) \ .master(local[*]) \ .getOrCreate() sc spark.sparkContext result sc.textFile(data.txt) \ .flatMap(lambda line: line.split( )) \ .map(lambda word: (word, 1)) \ .reduceByKey(lambda a, b: a b) \ .collect() for word, count in result: print(f{word}: {count})用Scala写就更短了val rdd sc.textFile(data.txt) val counts rdd.flatMap(_.split( )).map((_, 1)).reduceByKey(_ _) counts.saveAsTextFile(output/wordcount)有几个细节值得在复习时特别注意。第一textFile读的是文件也好、目录也好在分布式环境下会把文件按block切分成多个分区分区数决定了初始任务的并行度你可以在textFile方法的第二个参数里手动指定第二flatMap里的split如果遇到连续多个空格会拆出空字符串所以更严谨的写法是split(\\s)否则结果里会出现空字符串之类的小尾巴第三collect会把所有结果拉到Driver端如果结果集很大Driver内存容易爆正式场景一般用saveAsTextFile或take几个样本就够了。3.2 数据清洗项目里最常用的Spark套路说两个实际项目场景网约车订单日志和农产品价格数据。这两个场景里90%的清洗操作都是同一套模式读入数据、过滤脏数据、去重、处理空值、统一时间格式最后分组聚合。把这套流程练熟Spark SQL就已经算入门了。读CSV或JSON时用SparkSession的read方法远比RDD方便df spark.read.option(header, True) \ .option(inferSchema, True) \ .csv(orders.csv) clean_df df \ .filter(df[status].isin(完成, 取消)) \ .dropDuplicates([order_id]) \ .fillna({amount: 0.0}) \ .withColumn(order_date, to_date(create_time, yyyy-MM-dd HH:mm:ss))这里每步都可以单独拿出来当考点。filter里的isin是常见的多值过滤方式别只知道等于某个常量dropDuplicates按指定列去重和distinct按整行去重不一样fillna可以指定每列填充值金额为空的记录在后续聚合里不会出错withColumn加了order_date列把带时分秒的完整时间戳截断到天方便后续按天分组。分组聚合的阶段通常就是业务指标的计算时刻。比如网约车场景里按城市和日期统计订单量、总金额、均价农产品场景里按商品和月份统计平均价、最高价、最低价、环比变化。用Spark SQL的groupBy加agg可以一口气算出多个统计量result clean_df.groupBy(city_id, order_date) \ .agg( count(order_id).alias(order_cnt), sum(amount).alias(total_amount), avg(amount).alias(avg_amount) )如果面试官要求你写出TopN比如每个城市订单量最大的前十个司机那就需要窗口函数了。窗口函数在Spark SQL里是row_number() over (partition by ... order by ...)这部分许多人在课程里只是听说过但在真实的数据项目和面试题里出现频率特别高建议专门练习。3.3 数据倾斜项目里最常见的性能杀手数据清洗跑得慢的原因十有八九是数据倾斜。现象很典型Spark UI上某个Stage里绝大多数Task几十秒就结束了但有一两个Task挂了很久GC时间不断上升整个Stage的进度条卡在99%不动。原因通常是某个key的数据量特别大比如网约车订单里一个热门机场城市ID集中了大量订单groupBy时这个key所在的分区就要处理比其他分区多出几十倍的数据。处理思路分几层。最简单的是给shuffle调整分区数增大并行度但这治标不治本热点key还是会压在那一个分区上。更常用的套路是加盐salting给热点key加上随机前缀把它拆成多份先做一轮局部聚合去掉前缀后再做一轮全局聚合也就是两阶段聚合。以reduceByKey为例你可以先把结果做map(word (word _ random.nextInt(10), count))第一轮reduceByKey后变成map(kv (kv._1.split(_)(0), kv._2))第二轮reduceByKey得到最终结果。还有一类场景是join中的数据倾斜如果其中一个小表维度表可以广播那就是最优解Spark会直接把小表分发给所有Executor完全跳过shuffle。如果你的数据规模实在太大小表也广播不了可以尝试把热点key单独拆出来和对应的热点数据单独join后再union虽然代码复杂一些但往往能把任务时间从几小时压到十几分钟。注意数据倾斜不是一个配置能解决的必须先搞清楚倾斜发生在哪个Stage、哪个key上。先在Spark UI里观察各Task的输入数据量和执行时间再决定用什么方案别上来就调内存。4. Spark SQL专项日期、join与优化器4.1 DataFrame、Dataset与RDD的选择逻辑同样一段数据处理用RDD写和用DataFrame写性能可能差出一个数量级。原因在于DataFrame是带schema的结构化数据Spark能对整棵转换树做Catalyst优化比如调整join顺序、合并过滤条件、下推谓词到数据源。RDD则没有这些信息只能靠你自己手动优化。Dataset是DataFrame的类型安全版本它在编译期就能发现字段名错误但主要用在Scala API里Python里的DataFrame其实就是Dataset[Row]的别名。实际项目建议需要复杂灵活的数据变换、操作底层存储时用RDD业务逻辑清晰、以数据清洗分析为主时一律用DataFrame或者直接写Spark SQL因为SQL语法对后来维护的人更友好也更容易做自动优化。Spark 3.x之后还要提一个AQEAdaptive Query Execution。它在运行时动态调整shuffle分区数、自动处理数据倾斜、把sort merge join优化成broadcast join。这个机制的出现大大降低了新手调优的难度以前需要手动配置的很多参数现在Spark会自动做。但如果你的作业跑得很慢依然建议先在UI上确认AQE到底帮了什么忙别总觉得参数没调对。4.2 日期时间函数date_add、add_months这些高频操作“spark sql 日期加 年”“spark sql 日期加减”“spark sql 日期转存月份”这些搜索热词出现频率这么高说明日期处理真的是大家绕不过去的坎。日期函数在Spark SQL里有一批非常稳定的用法我把最高频的几个列成一张速查表函数作用例子date_add(date, n)日期加n天date_add(2024-01-31, 1) 2024-02-01date_sub(date, n)日期减n天date_sub(2024-03-01, 1) 2024-02-29add_months(date, n)日期加n个月add_months(2024-01-31, 1) 2024-02-29datediff(end, start)两个日期相差天数datediff(2024-03-01, 2024-02-01) 29months_between(end, start)相差月份数months_between(2024-03-01, 2024-01-01) 2last_day(date)该月最后一天last_day(2024-02-10) 2024-02-29date_format(date, fmt)按格式转字符串date_format(2024-01-05, yyyy-MM) 2024-01to_date(str, fmt)字符串转日期to_date(20240101, yyyyMMdd) 2024-01-01使用时有几个容易踩的坑得单独说。第一add_months处理月末日期有特殊语义如果源日期是1月31日加一个月会得到2月最后一天而不是3月初这在算账期、还款日时非常容易出错。第二字符串转日期时必须保证格式匹配to_date(2024-01-01, yyyy-MM-dd)和to_date(20240101, yyyyMMdd)互不兼容虽然Spark有时能自动识别标准格式但当你面对多种输入格式时最好先统一规范化再转换。第三date_format返回的是字符串如果你要再和日期做计算必须先转回日期类型否则类型不匹配直接报错。第四跨时区的时间戳处理要格外小心如果用to_timestamp解析带时区的时间最好明确目标时区否则不同的部署环境可能得出不同的结果。4.3 left outer join为什么只能广播右侧这个点能成为热搜词说明它很可能在面试里被问到过。大多数人知道Spark做join时可以广播小表但不知道广播的方向还受join类型限制。要理解这个问题先得把Broadcast Hash Join的机制说清楚。Broadcast Hash Join的过程是把参与join的两张表中较小的一张收集到Driver端构造成一个哈希表通过广播变量分发到所有Executor每个Executor上另一张表的每个分区逐行去探测这个哈希表。被广播的那张表在Spark源码里叫build side另一张表叫probe side。广播带来的收益是没有shuffle因为每个Executor都能本地拿到小表的完整哈希表。但left outer join的语义是“保留左表全部行右表没匹配上的补null”。如果把左表当成build side来构建哈希表让右表分区逐行去探测左表里那些在右表找不到匹配的行就会在探测过程中被遗漏无法完整保留只有让左表做probe side逐行去右表构建的哈希表里查查不到就补null才能完整输出左表的每一行。所以Spark在BroadcastHashJoinExec的源码里做了硬性限制LeftOuter、LeftSemi、LeftAnti这三种join类型build side只能选右表换言之只能广播右侧RightOuter则反过来只能广播左侧Inner Join两侧都能广播FullOuter虽然也能做但保留两边的语义处理起来更复杂实际项目中很少用广播方向去优化。注意即使你在SQL里显式写/* BROADCAST(左表) */Spark生成物理计划时还是会做build side合法性检查不合法就自动降级为Sort Merge Join并不会硬着头皮按照你的hint去执行。实际开发里最常见的正确姿势是左表大、右表小left join右表这个字典表右表能广播就广播整体跑起来会非常快。这个知识点不仅考试爱考在真实调优时也决定了你要不要给SQL加广播提示。5. 部署与调优从单机到集群的资源管理5.1 环境搭建与部署模式选择“spark集群搭建”“spark环境搭建及wordcount代码实现”“dgx spark 单机部署”这些热词出现很多次说明大家在入门时最头疼的就是环境。其实Spark的环境搭建本身不算复杂难的是搞清楚自己到底要哪种部署模式。最简单的是local模式直接在本地跑spark-shell --master local[*]星号表示使用本机所有可用核心数适合开发调试和课程作业。DGX这类GPU服务器上单机部署Spark本质上也是local模式加资源限制数据从本地文件系统读进去跑完测试后输出到本地目录它要的只是Spark运行环境不需要HDFS、YARN那些组件。如果你要在独立集群上跑就要用standalone模式先在一台节点执行start-master.sh再在计算节点执行start-worker.sh spark://主机名:7077然后通过8080端口看集群UI。这个模式的好处是Spark自带、部署简单缺点是没有统一的资源管理和队列机制适合学习和小规模自用。在生产环境里最常见的是YARN模式。你只需要在Spark的配置里设置HADOOP_CONF_DIR然后用spark-submit --master yarn提交作业YARN负责分配容器、管理资源。用这种模式时要注意一个关键点YARN的容器内存和Spark的executor内存是两个概念spark.executor.memory指定的是JVM堆内内存YARN分配的容器内存还包含堆外内存和开销如果你只调一个极容易出现“内存明明没满容器却被kill”的诡异情况。5.2 Spark内存模型与OOM排查思路Spark 2.x之后的内存模型是“统一内存管理”。Executor的JVM堆被分成三块预留内存Reserved Memory默认300MB不可动、用户内存User Memory存放用户自己创建的对象和Spark内存Spark Memory包括Storage Memory和Execution Memory。Storage用来缓存数据和RDDExecution用来做shuffle、join、聚合时的临时数据两块内存可以互相借用Storage有空闲时Execution可以用Execution用完再还。这个模型直接决定了为什么调参不能瞎调。spark.memory.fraction默认0.6意思是Spark Memory占整个JVM堆的比例spark.memory.storageFraction默认0.5表示Storage Memory初始占Spark Memory的比例。如果你的作业以shuffle和计算为主可以适当降低storageFraction让Execution有更多可用内存如果你有大量需要缓存的数据则反过来。但无论怎么调OOM的根源往往不是单个参数而是算子写法和数据分布出了问题。常见的OOM场景有三类。第一类Driver端OOM最常见的代码是collect一个超大的RDD回Driver比如几亿条记录直接拉到本地这种不需要调内存改代码就好用take抽样或者saveAsTextFile输出到分布式文件。第二类Executor端OOM跑shuffle时报Java heap space通常是有数据倾斜或者单个分区的数据量太大优先检查倾斜而不是无脑加内存。第三类被YARN直接杀掉日志里明确写“Container killed by YARN for exceeding memory limits”这说明你给Spark配的内存加堆外开销之后超过了容器内存需要适当把executor内存调小一点留出堆外空间或者给容器申请更大的内存。5.3 spark-submit参数速查与调参原则这里有一张我实际用下来最顺手的参数速查表考试和面试前可以快速扫一眼参数默认值用途说明--masterlocal[*]部署模式常见local、spark://host:7077、yarn--deploy-modeclientclient还是clustercluster模式下driver运行在集群里--executor-memory1g每个executor的JVM堆内存--driver-memory1gdriver端JVM堆内存--executor-cores1每个executor占用的核数--num-executors无固定executor数量YARN模式下常用spark.sql.shuffle.partitions200SQL shuffle默认分区数spark.default.parallelism随RDD分区变化RDD默认并行度spark.memory.fraction0.6Spark Memory占堆的比例spark.memory.storageFraction0.5Storage初始占Spark Memory的比例spark.serializerJavaSerializer改成KryoSerializer可提升序列化性能需要注册类spark.sql.adaptive.enabledtrue3.xAQE开关调参有一条我反复跟人强调的原则一次只改一个参数改完看作业在Spark UI上的表现别一次改一堆。很多新手一上来就把executor内存从2G改成16G、分区数从200改成2000结果作业没变快反而因为分区过多、调度开销增大而变慢。内存加得再大也不如把数据倾斜问题解决掉来得立竿见影。6. 高频报错与复习建议6.1 环境搭建期的典型报错我把这几年遇到过的报错按阶段整理了一份排查表先看环境搭建阶段。报错现象常见原因解决方法JAVA_HOME is not setSpark启动脚本找不到JDK在~/.bashrc里export JAVA_HOME和PATHspark-shell或spark-submit命令找不到SPARK_HOME未配置或未加入PATH设置SPARK_HOME并把$SPARK_HOME/bin加入PATH连接master失败Connection refusedstandalone模式下master没启动或地址写错先访问8080端口确认master状态再检查提交命令里的host和端口端口被占Address already in use8080或7077被占用改SPARK_MASTER_PORT、SPARK_MASTER_WEBUI_PORT或先杀掉占用进程Py4JNetworkError频繁出现Python环境和Spark版本不匹配或网络不稳定检查Python版本与PySpark版本兼容尽量用conda重建环境单机部署时还有个容易忽略的问题文件路径。如果你在本地磁盘上放了数据提交作业时要用file:///home/user/data.txt这种完整路径否则Spark默认去HDFS上找。很多人在单机环境第一次跑WordCount就报“输入路径不存在”结果只是路径前缀的问题不是代码问题。6.2 作业运行期的典型报错进入作业运行阶段这些报错要能一眼认出来。报错现象常见原因解决方法org.apache.spark.SparkException: Task not serializable闭包里引用了不可序列化的对象将类改成object把不必要字段改成局部变量或使用foreachPartition避免闭包捕获java.io.FileNotFoundException: output目录已存在saveAsTextFile要求输出目录不存在换输出路径或先删除已存在目录ExecutorLostFailure (executor lost)Executor所在节点OOM、心跳丢失或被外部kill看worker日志和Spark UI判断是内存问题还是网络问题Container killed by YARN for exceeding memory limitsexecutor内存加堆外超过YARN容器上限减小spark.executor.memory或增大yarn容器内存申请Py4JNetworkError: Connection refuseddriver与executor之间通信异常常见于网络或SparkContext被关闭检查网络连通性或者重启SparkContextjava.lang.OutOfMemoryError: Java heap space常见于单个分区数据过大或Driver端collect过大先查数据倾斜再考虑增大executor内存或改写collect这里我想专门强调一下Task not serializable。初学Spark的人第一次遇到这个报错往往一脸懵。报错本质是Spark的闭包会把你的函数序列化后传到Executor去执行如果你的函数里引用了某个类的实例而那个类没有实现Serializable接口就会失败。一个典型场景是你在main函数里new了一个helper对象然后在map里调用helper的某个方法。解决办法并不神秘把helper改成object单例对象或者把需要的字段提取成局部变量放到闭包里再不行就让类继承Serializable。6.3 我自己复习Spark时的几条实用建议最后分享一点个人经验尤其是马上要进考场或者面试的朋友可以参考这套复习路径。第一先跑通WordCount再把输出改成写入MySQL或者保存成Parquet文件。这一步能让你把数据读取、转换、聚合、输出、提交作业这五件事全部串起来比背十遍RDD概念都管用。第二学会看Spark UI。虽然课程考卷很少让你看UI但面试场景题里经常让你判断“某个Stage很慢是什么原因”这时候你对UI的熟悉程度直接决定你能不能答出关键点。第三做数据清洗练习时故意给自己出几道涉及日期、join、窗口函数的题比如“统计每天各个城市的订单量Top5商品”这道题覆盖了to_date、groupBy、row_number三个高频知识点一石三鸟。第四碰到报错不要急着搜答案先自己猜三个可能原因再去日志里验证。越是能安静下来读日志的人进步越快。第五把spark-submit常用参数抄在一张纸上每次提交作业都对着看一遍再执行时间长了这些参数自然就记住了而不是考前靠死记硬背。说实话Spark的知识点并没有多难难的是把它和实际场景联系起来。你已经看到这里了说明至少没被这类内容劝退接下来缺的只是亲手跑一遍。找一个真实的数据文件从WordCount开始慢慢往上加清洗、加聚合、加join等你把这条线走完再回头看这篇复习汇总你会发现每个知识点都有了一个落脚的场景。这就是我在这篇分享里最想传达的东西对于Spark这种工具型框架实践永远是比记忆更可靠的复习方法。