
这段时间在带内部新人培训正好讲到“3-4 Apache Spark基础”这一节。很多同学第一次接触 Spark 时容易把它理解成“一个更快版本的 Hadoop”但真正用起来才发现Spark 的定位不是简单替换而是一套完整的统一计算引擎批处理能跑流处理也能跑SQL、机器学习、图计算统统可以在同一套 API 上实现。这篇内容不是官方文档的搬运而是我这些年从入门到在生产环境调优一边踩坑一边总结出来的基础框架。如果你是刚接触 Spark 的开发者或者已经写过一些 RDD 代码但始终觉得概念比较散这一篇可以帮你把知识点串起来。我会从“Spark 到底解决了什么问题”开始讲到 RDD、DataFrame、Dataset 的区别再落到环境搭建、spark-submit 参数、常见报错排查最后顺手把启动时那行 log4j 日志刷屏的问题也处理掉。整个过程尽量说人话能直接照着操作。1. Apache Spark 是什么先搞清楚它解决的痛点1.1 从 MapReduce 的短板说起想理解 Spark绕不开 Hadoop MapReduce。MapReduce 这个模型本身很优雅把任务拆成 Map 和 Reduce 两个阶段中间通过 Key-Value 形式的数据流转。但它的短板同样明显——每个 Job 的中间结果都要落到磁盘下一个阶段再重新读上来。磁盘 I/O 在数据量大的时候就是瓶颈一个复杂的计算链路如果拆成好几个 MapReduce 串起来每一层都要写一次磁盘性能损耗非常可观。另外MapReduce 的 API 对开发者并不友好。你只是想做一次简单的分组求和也要写一堆 Mapper、Reducer、Driver 类编译、打包、提交整个流程非常重。交互式探索和分析场景根本跑不起来这也就催生了 Spark 这类“下一代”计算引擎的出现。Spark 的思路很直接把中间结果尽量留在内存里只有内存不够时才落盘。再加上它提供了一种比 MapReduce 更灵活的编程模型允许你把多个计算步骤直接串成一个 DAG有向无环图由引擎统一优化执行。这样一来同样一个多阶段任务Spark 可能只需要读一次数据MapReduce 却要反复读写磁盘差距就拉开了。1.2 Spark 的核心卖点内存计算与统一引擎很多人把“内存计算”四个字背得很熟却不太清楚它到底快在哪。举个生活中的例子想象你是一个厨师MapReduce 的做法是每做完一道菜就把所有材料搬回仓库下一道菜再从仓库搬出来Spark 的做法是材料就放在操作台上做完一道菜直接顺手用剩下的材料做下一道。数据在内存里流动自然比在磁盘上反复搬运快得多。但 Spark 能流行起来不仅仅是因为快还因为它的“统一”统一 API一套 RDD/DataFrame/Dataset 的代码既能写批处理也能写流处理。内置 SQLSpark SQL 让你直接写 SQL 查询表格数据不需要额外学一套新语法。生态丰富MLlib 做机器学习、GraphX 做图计算、Structured Streaming 做流处理都挂在同一个核心之上。所以学 Spark 基础不建议一上来就抱着源码啃。你先要建立一个全局观它是一个计算引擎不是一个存储系统数据可以来自 HDFS、S3、本地文件或者各种数据库它负责的是“算”不是“存”。这个定位想清楚后后面很多概念就顺了。2. 核心概念拆解RDD、DataFrame、Dataset 到底怎么选2.1 RDD最底层的抽象RDD 的全称是 Resilient Distributed Dataset中文常翻译成“弹性分布式数据集”。我习惯把它理解成“一个分布在多台机器上、可以并行操作、并且出错后能重新计算的集合”。它有三个关键特点分区Partition数据被切成若干分区每个分区在不同的 Executor 上被并行处理。不可变ImmutableRDD 一旦创建就不能改只能通过 transformation 生成新的 RDD。惰性求值Lazy Evaluationtransformation 只是记录计算过程不真正执行只有遇到 action比如 collect、count、save时才会真正开始算。初学阶段最常见的坑就是把 transformation 当成“已经执行完”。比如有些人写完rdd.map(...)马上打印这个 rdd 去看结果发现里面还是空的就是因为 map 只是定义了“将来要做什么”并没有跑。这是 Spark 和普通集合 API 最大的区别也是惰性求值的核心价值引擎可以看到整个计算链条再决定怎么优化。不过在实际工作中RDD 用得越来越少。原因也很简单——它没有 schema 信息引擎没法针对数据结构做优化而且框架代码里浮现大量map(x ...)、reduceByKey(...)的 Lambda 表达式可读性和维护性都一般。2.2 DataFrame / Dataset面向结构化数据的更高层 APIDataFrame 解决的就是 RDD“没有 schema”的痛点。你可以把它理解成一张分布式表每一行是一个对象每一列有明确的类型。由于引擎知道了列的类型就能做很多自动优化比如谓词下推、列裁剪、二进制存储。Dataset 是强类型版的 DataFrame。在 Scala/Java 里你可以定义一个case class Person(name: String, age: Int)然后得到一个Dataset[Person]编译时就能发现字段写错的问题。在 Python 里DataFrame 和 Dataset 的边界比较模糊你实际接触到的更多是 DataFrame。三者的关系可以粗略画成一条链DataFrame 是Dataset[Row]的特例RDD 位于最底层而 DataFrame/Dataset 在 RDD 之上加入 schema 和执行优化。使用 Spark SQL 时你写的 SQL 会被解析成逻辑计划最终转换成 RDD 上的物理执行计划但这个过程不需要你手动介入。2.3 实际选型建议给新人的建议其实很简单能用 DataFrame / Spark SQL 写就尽量不写 RDD。只有在需要精细控制底层数据分布或者使用第三方 RDD API 时才考虑直接操作 RDD。数据源来自 Hive、Parquet、JSON、JDBC 时优先建表或读成 DataFrame后面做过滤、聚合都会自动优化。维度RDDDataFrame / DatasetSchema 信息无有自动优化很少Catalyst 优化器类型安全弱Dataset 强类型适合场景底层自定义逻辑日常数据加工、SQL分析代码可读性一般较高我自己带项目时有个原则除非某个功能只有 RDD API 能实现否则一律用 DataFrame 起步。这样团队协作时沟通成本低性能也不会差。3. 作业执行原理从提交到出结果中间发生了什么3.1 Driver、Executor、Cluster Manager 的分工Spark 应用跑起来后进程层面有这么几个角色Driver你的 main 函数在这里执行负责解析代码、生成 DAG、把任务切分成 task、调度到 Executor 上跑最后汇总结果。Executor真正干活的进程运行在 Worker 节点上负责执行 task并把计算结果返回给 Driver。Cluster Manager负责任务资源分配可以是 Standalone、YARN 或 Kubernetes。三者的关系有点像开发团队Driver 是项目经理负责拆任务、派活、收工Executor 是一线开发真正出代码Cluster Manager 是 HR负责招聘和分配办公位。很多人刚开始看 Spark UI 时看到一堆 Executor 状态不太理解其实就是这么个协作关系。任务执行时Driver 会把用户代码转换成一个逻辑计划再转换成物理执行计划最终生成一个 DAG 图。DAG 会按照宽依赖切成若干个 Stage每个 Stage 内部再拆成一组可以并行执行的 Task这些 Task 分发到 Executor 上运行。这个过程在 Spark UI 的 Jobs / Stages 标签页里能看得非常清楚。3.2 Stage 划分、窄依赖与宽依赖Stage 的划分依据是依赖关系窄依赖父 RDD 的每个分区最多被一个子 RDD 分区使用。比如map、filter这类操作不需要跨节点 shuffle所在 Stage 内部可以直接 pipeline 执行。宽依赖父 RDD 的多个分区需要被同一个子 RDD 分区使用典型操作是groupByKey、reduceByKey、join。这种依赖通常会引起 shuffle是 Stage 划分的边界。可以这么理解窄依赖像流水线上每个工人只处理自己手头的零件宽依赖像把所有零件先倒进一个分拣中心再按某种规则重新分组发给下一批工人。分拣这个动作本身很费事所以一个作业里 shuffle 越多性能越容易出问题。3.3 为什么 shuffle 是性能杀手Shuffle 的本质是让数据跨节点重新分布。比如reduceByKey需要把相同 key 的数据拉到同一个节点再聚合。这个过程中会发生数据序列化、网络传输、磁盘读写、内存缓冲随便哪一环慢了整个作业都会拖慢。我在调优时经常会先看 Spark UI 里的 Shuffle Read / Shuffle Write 指标。如果某个 Stage 的 Shuffle Read 数据量特别大优先考虑是否有办法减少 shuffle比如用reduceByKey替代groupByKey前者可以在 map 端先做一次本地聚合或者调整分区数让数据分布更均匀。有一点要注意不是所有 join 都一定触发完整 shuffle。如果一张表足够小Spark 会自动选择 Broadcast Join把小表广播到每个 Executor不产生 shuffle性能提升非常明显。这也是为什么后续做优化时先看执行计划比盲目调参数更有效。4. 本地环境搭建与第一个 Spark 程序实操记录4.1 安装目录与 spark-shell 启动这里以 Spark 3.3 版本为例。下载二进制包后解压目录里常见的内容有bin、sbin、conf、jars等。bin/spark-shell是交互式 Scala 壳bin/pyspark是 Python 交互环境bin/spark-submit用来提交打包好的作业。启动本地模式时直接执行cd /opt/spark bin/spark-shell --master local[2]local[2]表示在本地用 2 个线程模拟 Executor。这个模式非常适合学习和调试不需要额外搭集群。启动过程中常见会刷出一行提示Using Sparks default log4j profile: org/apache/spark/log4j-defaults.properties这行不是报错只是告诉你 Spark 没找到自定义的log4j2.properties所以使用了默认日志配置文件。对初学者来说默认配置最烦人的是 INFO 级别的日志太多把真正有价值的 WARN 和 ERROR 信息都淹没了。4.2 用 spark-shell 跑一个 WordCount在spark-shell里可以直接写 Scala 代码。下面是最经典的 WordCountval rdd sc.textFile(file:///tmp/words.txt) val counts rdd .flatMap(_.split( )) .map(word (word, 1)) .reduceByKey(_ _) counts.collect().foreach(println)注意flatMap和map都只是定义操作直到你调用collect()这个 action 时Spark 才会真正开始调度执行。执行完成后控制台会输出每个单词的出现次数。如果数据在本地路径要写成file:///tmp/words.txt如果文件在 HDFS 上就写成hdfs://namenode:8020/tmp/words.txt。新手最容易在这里踩坑写成hdfs:///tmp/words.txt但本地没有 NameNode结果一直报错找不到文件。4.3 顺手处理默认 log4j 日志刷屏问题第一次跑任务时控制台会刷大量 INFO 日志看起来很吓人。解决办法很简单在conf目录下复制一份 log4j 配置文件然后调整日志级别。在 Spark 3.x 中默认文件是log4j2.propertiescd /opt/spark/conf cp log4j2.properties.template log4j2.properties编辑log4j2.properties找到类似rootLogger.level info改成rootLogger.level warn如果你想单独把 Spark 类的日志调成 WARN而保留其他类的 INFO可以这样加logger.spark.name org.apache.spark logger.spark.level WARN改完重启 spark-shell你会发现终端清爽多了。这一步看起来和“Spark 基础”没什么关系但在日常开发中非常实用。很多新人一上来看到满屏日志以为是程序出错了其实只是默认日志级别的正常输出。5. 提交作业到集群spark-submit 的常用姿势5.1 核心参数说明本地跑通了下一步就是把作业提交到集群。最常用的命令就是spark-submit。我一般会维护一个模板bin/spark-submit \ --master yarn \ --deploy-mode cluster \ --class com.example.WordCount \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 20 \ --driver-memory 2g \ wordcount.jar \ /input/words.txt /output/counts几个参数的作用要搞清楚--master指定集群管理器yarn、spark://、k8s:// 等。--deploy-modedriver 运行在客户端本地还是集群内部client 或 cluster。--executor-memory每个 Executor 申请多少内存。--executor-cores每个 Executor 使用多少个 CPU 核心。--num-executors一共启动多少个 Executor。--driver-memoryDriver 进程内存。这些参数不是随便填的。Executor 内存设太大会导致单台机器上只放得下一两个 Executor资源碎片化设太小shuffle 数据又要频繁落盘。通常我会先看数据集大小和计算复杂度再结合每台物理机的内存做分配。5.2 client 与 cluster 模式对比这里的核心区别是 Driver 进程跑在哪。client 模式Driver 跑在提交作业的客户端机器上适合调试因为日志能直接打印在终端。缺点是一旦客户端断开作业可能挂掉。cluster 模式Driver 被调度到集群内部跑适合生产环境日志要到 YARN 的容器里查看。即使你本地关机作业也照常执行。如果只是学习用 client 方便如果是定时调度或线上任务尽量用 cluster。否则某天你关掉笔记本电脑所有作业跟着断掉哭都来不及。5.3 内存和并行度配置经验关于内存有条经验值得记住--executor-memory不完全是堆内内存。Spark 的内存分为执行内存execution、存储内存storage和保留内存reserved默认还有一个比例参数spark.memory.fraction控制执行与存储共用部分的比例。所以把 executor-memory 调大不一定代表可用的堆内内存就大还要看spark.memory.offHeap.enabled和spark.memory.offHeap.size等参数。并行度方面新人最容易犯的错是根本不设置分区数。默认分区可能很小导致几百个核只有十几个 task 在跑资源大量空转。我通常在读取大表或做 shuffle 时显式指定分区数比如repartition(200)或者设置spark.sql.shuffle.partitions200。这个数值不固定一般以 executor 总核心数的 2 到 3 倍作为起点再根据任务耗时逐步调整。6. 新手最容易踩的坑与排查思路6.1 内存溢出OOMOOM 是 Spark 新手遇到的第一个“大魔王”。常见原因有三类数据加载过大、shuffle 数据量大、单个 partition 数据不均衡。出现 OOM 时不要急着把内存调大先看看 Spark UI 里的内存趋势和 task 耗时。如果某个 Stage 的 task 长时间不结束且数据倾斜明显优先考虑加宽 key 的分布比如加盐、两阶段聚合如果只是单纯数据量大可以先过滤掉无用的列和行再考虑增加 Executor 数量而不是单 Executor 内存。一上来就盲目把 executor-memory 调到 32g通常只会让 GC 更频繁反而更慢。6.2 数据倾斜的典型表现数据倾斜的表现很典型整个作业大部分 task 很快结束只剩一两个 task 一直卡在那里跑。原因通常是某个 key 的数据量特别大导致所有数据都堆到了同一个 partition。排查时先用 SQL 或 groupBy 统计一下 key 的分布。如果确认是少数 key 倾斜可以给 key 加随机前缀把数据打散到多个分区聚合后再去前缀做一次最终聚合。对于 join 倾斜可以先把小表广播出去避免因为热点 key 引发的 shuffle。这个优化技巧在面试里也常被问到但理解原理比背结论更重要数据倾斜本质上是“并行度失效”单个分区的处理时间决定了整个 Stage 的耗时。6.3 学会看日志和 Spark UI遇到问题时第一反应不是猜而是去看日志和 UI。Spark UI 的 Jobs、Stages、Executors 三个页面基本能回答绝大多数问题Jobs 页面可以看到每个 action 对应的 Job 总数和耗时。Stages 页面定位慢 Stage看是输入量大还是 shuffle 耗时长。Executors 页面看每个 Executor 的内存、GC 时间和读写数据量是否均衡。日志方面如果你已经按照前面 4.3 小节把 log4j 默认级别调成 WARN启动时那行Using Sparks default log4j profile就不会再频繁打扰你。真出问题需要看细节时再临时用--conf spark.log.levelINFO或者手动把配置改回来。我在实际带人时经常会说一句话Spark 基础要学的不是“把命令跑通”而是学会在任务跑不动的时候能通过日志和 UI 快速判断问题出在哪个环节。这个能力一旦有了后面学调优、学流计算都会事半功倍。如果只是按文档敲一遍代码过两周就忘了如果带着“每个概念解决什么真实问题”的思路去学你才能真正把 Spark 用起来。最后分享一个小习惯每拿到一套新的 Spark 环境我会先花五分钟做一次spark-shell启动、读一个本地文件、跑一个 count再打开 UI 看看执行计划。这五分钟能帮你排除掉大部分环境问题也是我最推荐的入门热身动作。