
简介这份资源面向零基础到进阶的大数据学习者围绕2020年9月发布的Spark 3.0.1稳定版展开覆盖从环境搭建到性能调优的完整学习路径适合希望系统掌握Spark核心组件与实战技巧的开发者。压缩包共244个文件以217张png截图、10份md笔记、9个zip示例包为主辅以7个scala源码和1个json配置整体约86.9MB截图与笔记便于对照理解源码与压缩包可直接运行验证。内容涵盖SparkCore、SparkSQL、SparkStreaming、StructuredStreaming、多语言开发、综合案例及3.0新特性等9个章节笔记按天拆分配合案例代码可帮助读者快速搭建实验环境、理解RDD与DataFrame编程模型、掌握流处理与调优思路。目前已有609人学习下载适合作为入门到精通的系统化参考资料。1. Spark 3.0 入门到精通8 天代码笔记到底能解决什么很多人第一次接触大数据卡住的地方不是“不会写代码”而是环境跑不起来、数据读不进来、任务一提交就报错。这份「大数据入门 Spark 3.0 入门到精通 1-8day 代码-笔记」本质上是一套按天推进的学习路径从 Spark 安装、RDD 与 DataFrame 基础到 Spark SQL、读取 JSON、内存调优和集群搭建每一天都配了可运行的示例代码和笔记。它适合两类人一是刚转大数据方向、需要一套能照着敲的入门材料二是已经会写 SQL 但没系统跑过 Spark 任务、想补齐分布式执行认知的开发者。下面我不复述笔记内容而是按一线落地的顺序把这条路径拆成能复现的步骤、参数和踩坑点。2. 环境与第一个 Spark 任务本地模式怎么跑通2.1 为什么先跑 local 模式而不是直接上集群新手最容易翻车的地方是一上来就照着「Spark 集群搭建」的教程配三台机器结果卡在 SSH 免密、时间同步、防火墙这些和 Spark 本身无关的环节上三天过去连一行代码都没跑。我的建议是先用 local 模式把代码逻辑跑通再迁移到 standalone 或 YARN。local 模式把 Driver 和 Executor 放在同一个 JVM 进程里没有网络通信和资源调度报错信息也干净得多。Spark 3.0 对 JDK 和 Scala 版本有要求这是第一个硬门槛。JDK 必须是 8 或 11Scala 2.12 是官方编译版本。用 JDK 17 会直接抛IllegalAccessError因为模块化之后反射访问被限制了。这一点在笔记里通常会写但很多人装环境时用的是系统自带的最新 JDK跑起来才发现。# 检查 JDK 版本必须是 8 或 11 java -version # 解压 Spark 3.0 安装包版本号按你实际下载的填 tar -zxvf spark-3.0.0-bin-hadoop3.2.tgz -C /opt/ # 配置环境变量写入 ~/.bashrc export SPARK_HOME/opt/spark-3.0.0-bin-hadoop3.2 export PATH$SPARK_HOME/bin:$PATH # 让配置生效 source ~/.bashrc # 启动 Spark 自带的交互式 shell验证安装 spark-shellSPARK_HOME指向解压目录PATH里加上bin是为了能直接调用spark-shell、spark-submit这些命令。spark-shell启动后会打印 Spark 版本和 Scala 版本看到Spark context available as sc就说明本地模式通了。如果卡在启动阶段先看$SPARK_HOME/logs下的日志九成是 JDK 版本或内存不足。2.2 用 spark-submit 提交第一个可复现任务spark-shell适合交互调试但真正的工作流是写脚本然后用spark-submit提交。下面这段代码统计一个文本文件里的词频是 Spark 的「Hello World」但我会把参数写全方便你直接改。# wordcount.py from pyspark import SparkContext, SparkConf # 创建配置setAppName 是任务名setMaster 指定运行模式 conf SparkConf().setAppName(WordCount).setMaster(local[2]) sc SparkContext(confconf) # 读取本地文件textFile 支持本地路径和 HDFS 路径 lines sc.textFile(file:///opt/data/words.txt) # flatMap 把每行拆成单词map 转成 (word, 1)reduceByKey 聚合 counts lines.flatMap(lambda line: line.split( )) \ .map(lambda word: (word, 1)) \ .reduceByKey(lambda a, b: a b) # collect 把结果拉回 Driver 端打印生产环境慎用 for word, count in counts.collect(): print(f{word}: {count}) sc.stop()setMaster(local[2])里的 2 表示用两个线程模拟并行本地调试够用。textFile的路径必须带file://前缀否则 Spark 会去 HDFS 找。collect()会把所有结果拉到 Driver 内存数据量大时直接 OOM这里只是演示。提交命令spark-submit \ --master local[2] \ --name WordCount \ /opt/scripts/wordcount.py--master和代码里的setMaster重复时命令行优先级更高。生产环境把local[2]换成yarn或spark://host:7077其余代码不用动这就是 Spark 提交接口统一的好处。3. RDD 与 DataFrame两套 API 怎么选、怎么转3.1 RDD 的惰性求值与宽窄依赖RDD 是 Spark 最早的计算抽象核心就两点不可变、惰性求值。你写的map、filter不会立刻执行只有遇到collect、count、saveAsTextFile这类行动算子才触发。这个设计让 Spark 能在真正执行前优化整个 DAG。理解宽窄依赖是排查性能问题的关键窄依赖map、filter每个父分区只对应一个子分区不需要 shuffle宽依赖groupByKey、reduceByKey需要把数据按 key 重新分发到不同节点触发网络传输和磁盘落盘。血泪经验是能用reduceByKey就别用groupByKey。reduceByKey会在 map 端先做本地聚合再 shuffle 聚合结果groupByKey把所有原始数据都 shuffle 过去数据量大一个数量级。下面这段代码演示两者的差别。# 假设 rdd 是 (key, value) 格式 rdd sc.parallelize([(a, 1), (a, 2), (b, 3), (b, 4)]) # 推荐map 端预聚合shuffle 数据量小 reduced rdd.reduceByKey(lambda x, y: x y) # 不推荐全量 shuffle容易 OOM grouped rdd.groupByKey().mapValues(sum) # 两者结果一样但执行计划完全不同 print(reduced.collect()) # [(a, 3), (b, 7)]reduceByKey的聚合函数要求满足交换律和结合律否则并行计算结果可能不对。groupByKey返回的是迭代器mapValues(sum)会把它转成列表再求和内存压力全在 Executor 上。3.2 DataFrame 与 Spark SQL结构化数据的正确打开方式Spark 3.0 里 DataFrame 是主流底层还是 RDD但带了 Schema 信息Catalyst 优化器能做的事比手写 RDD 多得多。读取 JSON 是热搜里高频出现的场景spark.read.json会自动推断 Schema但推断有代价它要扫一遍全量数据。生产环境建议显式指定 Schema省掉这次扫描。from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType spark SparkSession.builder \ .appName(JsonRead) \ .master(local[2]) \ .getOrCreate() # 显式定义 Schema避免全量扫描推断 schema StructType([ StructField(name, StringType(), True), StructField(age, IntegerType(), True), StructField(city, StringType(), True) ]) # 读取 JSON指定 schema df spark.read.schema(schema).json(file:///opt/data/users.json) # 创建临时视图用 SQL 查询 df.createOrReplaceTempView(users) result spark.sql(SELECT city, COUNT(*) AS cnt FROM users GROUP BY city) result.show() spark.stop()StructField的第三个参数表示是否允许 null设成False时遇到空值会报错能帮你尽早发现脏数据。createOrReplaceTempView注册的视图只在当前 Session 有效跨 Session 要用createGlobalTempView。show()默认只显示 20 行调试够用别拿它当导出工具。4. 避坑与排查新手最常翻车的 5 个场景4.1 现象任务卡在Job aborted due to stage failure日志里全是 shuffle 失败原因通常是 Executor 内存不够shuffle 阶段数据落盘时 OOM。Spark 3.0 默认 shuffle 分区数是 200小数据量下每个分区任务很轻大数据量下单个分区可能几个 G。解决方式是调spark.sql.shuffle.partitions按数据量估算一般设成 Executor 核数的 2 到 3 倍。spark-submit \ --conf spark.sql.shuffle.partitions400 \ --conf spark.executor.memory4g \ --conf spark.executor.cores2 \ your_script.py4.2 现象java.lang.OutOfMemoryError: Java heap space出现在 Driver 端原因是你用了collect()或toPandas()把大结果集拉回本地。Driver 默认内存只有 1G几百万行数据直接撑爆。解决方式是改用write落盘或者调大spark.driver.memory。但调大内存只是缓解根本办法是别把分布式数据往单机拉。4.3 现象读取 JSON 时字段全是 nullSchema 推断失败原因有两种一是 JSON 文件里有跨行记录spark.read.json默认按行解析遇到多行 JSON 会解析失败二是字段类型混杂比如 age 字段有的行是数字有的行是字符串推断成 StringType 后数字查询失效。解决方式是加multiLineTrue参数或者显式指定 Schema。# 处理跨行 JSON df spark.read.option(multiLine, true).json(file:///opt/data/multi.json)4.4 现象spark-shell启动报NoClassDefFoundError或ClassNotFoundException原因是 Spark 的 jar 包和 Hadoop 版本不匹配。Spark 3.0 预编译版有hadoop2.7和hadoop3.2两个版本下载时选错就会在读写 HDFS 时报错。解决方式是重新下载对应版本或者用--jars手动指定 Hadoop 客户端 jar。4.5 现象本地模式跑得好好的提交到集群后找不到文件原因是路径问题。本地模式file:///指向本机文件系统集群模式下每个 Executor 在不同的机器上file:///指向的是 Executor 所在机器的本地路径文件根本不存在。解决方式是把数据放到 HDFS 或对象存储代码里用hdfs://或s3a://路径。这个坑几乎每个人都会踩一次记住「集群模式下没有本地文件」这句话能省半天排查时间。5. 内存调优与集群提交从能跑到跑得稳5.1 Spark 内存模型Execution 与 Storage 怎么分Spark 3.0 的 Executor 内存分三块Executionshuffle、join、sort 用、Storage缓存 RDD、广播变量用、Other用户代码和元数据。统一内存管理下Execution 和 Storage 可以互相借用但 Execution 优先级更高Storage 被借走的部分如果还要用就得重算。调优的核心是让这两块别互相抢。关键参数是spark.memory.fraction默认 0.6表示堆内存里 60% 给 Execution 和 Storage。剩下 40% 留给用户代码和元数据。如果你的任务 shuffle 特别重可以调到 0.7 或 0.8如果缓存了大量 RDD也要保证 Storage 有足够空间。spark-submit \ --conf spark.memory.fraction0.7 \ --conf spark.memory.storageFraction0.3 \ --conf spark.executor.memory8g \ --conf spark.executor.memoryOverhead2g \ your_script.pyspark.memory.storageFraction默认 0.5表示统一内存里 Storage 占一半。调小它意味着 Storage 更容易被 Execution 挤占缓存数据可能被驱逐。spark.executor.memoryOverhead是堆外内存用于 JVM 自身开销和 Python 进程PySpark 任务建议设成 executor 内存的 10% 到 20%。5.2 集群提交standalone 与 YARN 的取舍standalone 是 Spark 自带的调度器部署简单适合专用集群。YARN 是 Hadoop 生态的资源调度器适合和 MapReduce、Hive 共享集群。选哪个取决于你现有的基础设施如果已经有 Hadoop 集群直接用 YARN别另起炉灶如果是全新环境且只跑 Sparkstandalone 更轻量。standalone 提交命令spark-submit \ --master spark://master-host:7077 \ --deploy-mode cluster \ --executor-memory 4g \ --total-executor-cores 8 \ your_script.pyYARN 提交命令spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --num-executors 4 \ --executor-cores 2 \ your_script.py--deploy-mode cluster表示 Driver 跑在集群里客户端断开也不影响任务client模式 Driver 在本地适合调试但客户端一关任务就断。--num-executors和--total-executor-cores是两种调度器的不同表达YARN 用前者standalone 用后者。提交后去 Web UI 看任务状态standalone 默认 8080 端口YARN 看 ResourceManager 的 8088 端口。6. 进阶技巧用执行计划定位性能瓶颈6.1 看懂explain()输出找到多余的 shuffle很多人调优靠猜改参数试来试去。其实 Spark 的explain()会把逻辑计划和物理计划都打出来哪里 shuffle、哪里过滤下推一目了然。下面这段代码对比两种写法的执行计划。# 写法一先 join 再 filter df1 spark.read.json(file:///opt/data/orders.json) df2 spark.read.json(file:///opt/data/users.json) result1 df1.join(df2, user_id).filter(df1.amount 100) result1.explain() # 写法二先 filter 再 join df1_filtered df1.filter(df1.amount 100) result2 df1_filtered.join(df2, user_id) result2.explain()explain()输出里看到Exchange就是 shuffle看到Filter在Join之前就是谓词下推生效。写法二把过滤提前join 的数据量变小shuffle 数据量也跟着小。这个优化 Catalyst 有时能自动做但涉及多表 join 时不一定手动调整更保险。6.2 广播变量与小表 join当一张表小到能放进 Executor 内存时用广播 join 可以完全避免 shuffle。Spark 3.0 默认广播阈值是 10MB可以通过spark.sql.autoBroadcastJoinThreshold调整。超过阈值的大表广播会导致 Driver OOM所以别盲目调大。from pyspark.sql.functions import broadcast # 显式广播小表 result big_df.join(broadcast(small_df), id)broadcast()提示优化器把 small_df 分发到每个 Executorjoin 在本地完成。判断标准是 small_df 压缩后的大小不是行数。一个 100 万行但只有两列的表可能只有几 MB照样能广播。6.3 数据倾斜加盐与两阶段聚合数据倾斜是分布式计算的玄学问题大部分 task 几秒跑完个别 task 跑几小时。原因是某个 key 的数据量远超其他 key。解决办法是给 key 加随机前缀打散聚合后再去掉前缀二次聚合。from pyspark.sql.functions import rand, concat, lit, col # 加盐给倾斜 key 加 0-9 的随机前缀 salted df.withColumn(salted_key, concat(col(key), lit(_), (rand() * 10).cast(int))) # 第一次聚合按加盐后的 key 聚合 partial salted.groupBy(salted_key).sum(value) # 去掉盐第二次聚合 result partial.withColumn(key, split(col(salted_key), _)[0]) \ .groupBy(key).sum(sum(value))加盐的代价是多一次 shuffle所以只对确认倾斜的 key 用别全量加。判断倾斜的方法是看 Spark UI 里 stage 的 task 耗时分布如果 max 是 median 的几十倍基本就是倾斜了。我自己的习惯是每次写完一个任务先explain()看一眼执行计划再跑一个小数据集验证逻辑最后才上全量数据。这个顺序能挡掉八成低级错误。希望帮到你。本文还有配套的精品资源点击获取