简介面向大数据初学者和Spark入门用户这套基于Spark3.0.1的代码与笔记按1-8天的学习路径编排覆盖环境搭建、SparkCore、SparkStreaming、SparkSQL、StructuredStreaming、综合案例、多语言开发、3.0新特性及性能调优共九个章节可用于系统入门、课堂配套和日常复习。压缩包共244个文件约86.9MB以217张png截图、10个md笔记和7个scala代码为主体另有9个zip子包和1个json配置截图帮助你对照操作与运行结果md整理各日学习要点scala为可直接运行的练习代码zip则容纳了配套资源。已有609人学习下载对正在入门Spark并希望构建完整知识体系的读者是一份组织清晰、涵盖核心组件的实用资料。1. 拿到这份 spark3.0 学习包之后先别急着复制粘贴解压“大数据入门spark3.0入门到精通 1-8day 代码-笔记.zip”之后多数人第一反应是打开 day1 的笔记复制代码点运行然后卡在第一个报错上。这个包里真正值钱的不是那些能跑的代码而是按天规划的路线环境搭建、RDD、DataFrame、SQL、结构化流、调优八天正好覆盖大数据从“看懂”到“能写”的跨度。适合有 Python 或 SQL 基础、想往数仓或实时计算方向走的人。这里我换一种学法不按笔记从头读到尾而是先把 Spark 3.0 的运行机制立住照着下面的步骤搭好环境、跑通第一个作业再倒回去翻那八天笔记你会发现读代码的速度完全不一样。2. Spark 3.0 入门必须补的三块地基RDD、DAG 与 AQE任何一份 Spark 入门笔记前几节都会讲 RDD、DAG、shuffle 这些名词。但如果只是扫一眼就往下走后面读代码会越来越吃力。这里把三块最影响你理解代码的地基先讲清楚它们是后面 8 天路线里反复出现的东西。2.1 先搞清 RDD、DataFrame、Spark SQL 的关系再决定先学哪个RDD 是 Spark 最底层的抽象它把数据切成多个分区每个分区由一个 Task 处理。DataFrame 在 RDD 之上加了一层 Schema也就是每一列叫什么、什么类型这让 Spark 能在执行前优化你的代码。Spark SQL 则把 DataFrame 转成 SQL 语句来跑适合数仓场景。入门阶段最常见的误区是花大量时间背 RDD 算子map、flatMap、reduceByKey 背了一堆结果真去做业务时发现大家都在写 Spark SQL。我的建议是RDD 的常用算子过一遍理解分布式计算是怎么回事主力练习放在 DataFrame 和 Spark SQL 上。下面这段代码用两种写法做同一件事你对比一下差别。from pyspark.sql import SparkSession spark SparkSession.builder.appName(compare_rdd_df).master(local[2]).getOrCreate() # 准备一份本地日志文件每行是一行访问日志 log_rdd spark.sparkContext.textFile(file:///tmp/access.log) # RDD 写法明确告诉引擎每一步怎么做 ip_counts_rdd log_rdd.map(lambda line: line.split( )[0]) \ .map(lambda ip: (ip, 1)) \ .reduceByKey(lambda a, b: a b) # DataFrame 写法只告诉引擎要什么结果 log_df spark.read.text(file:///tmp/access.log) ip_counts_df log_df.selectExpr(split(value, )[0] as ip) \ .groupBy(ip).count() \ .orderBy(count, ascendingFalse) ip_counts_df.show()逻辑说明RDD 写法里每一步转换都是手动指定的map 取 IP、map 组装键值对、reduceByKey 做聚合引擎不知道你想干什么只能照做。DataFrame 写法里引擎看到的是“按 IP 分组统计数量”Catalyst 优化器会帮你做谓词下推、列裁剪。参数说明local[2]表示用本地两个线程模拟两个分区适合学习阶段验证逻辑真实提交到集群时要换成具体资源参数。2.2 一次作业从提交到出结果DAG、Stage、Task 是怎么流转的理解作业执行流程是后面看日志、做调优的前提。一次spark-submit提交后Driver 端会构建一个 DAG也就是执行计划图。DAG 里的每一步转换会根据宽依赖被切成多个 Stage宽依赖最典型的就是 shuffle比如groupBy、reduceByKey、join。Stage 再拆成多个 Task分发给 Executor 并行执行。看执行计划不是看玄学explain(true)是入门阶段必须养成的习惯它会把物理计划全部打出来。from pyspark.sql import SparkSession spark SparkSession.builder.appName(explain_demo).master(local[2]).getOrCreate() df spark.read.option(header, True).csv(file:///tmp/data.csv) df.filter(age 0).groupBy(city).count().explain(true)逻辑说明explain(true)会输出解析逻辑计划、优化后逻辑计划、物理计划三部分。入门阶段你只需要看最后一段物理计划重点找两个地方一是Exchange它代表 shuffleshuffle 越多作业越慢二是Scan它代表读数据的方式如果数据量大而且没走分区裁剪就得怀疑写查询的方式。参数说明true表示打印详细计划不带参数时只打印物理计划摘要。这个命令在 pyspark 和 spark-sql 里都能用是排查“为什么这么慢”的第一现场。理解 DAG 还有一个实际用途看日志里 Stage 的数量。同一个需求用groupBy是一个 Stage如果你用groupByKey或反复repartitionStage 数量会变多shuffle 次数也变多跑得慢就不奇怪了。2.3 直接学 3.x 的理由AQE 和动态分区裁剪让新手少背几十个参数网上大量教程还停留在 Spark 2.x 的思路动不动就让你调十几个参数。Spark 3.0 引入的 AQE自适应查询执行和大规模应用的动态分区裁剪把很多手动调优变成了自动行为。AQE 会在作业执行过程中根据已完成 Task 的统计结果动态调整后续计划它主要做三件事自动合并小的 shuffle 分区、自动切换 join 策略、自动处理数据倾斜。对入门者来说最直观的变化是以前处理数据倾斜要加盐、要手动判断大 key现在只要把 AQE 打开很多场景它会自动拆分倾斜分区。先确认你当前环境里 AQE 到底开没开再决定要不要手动配参数。spark-submit --versionfrom pyspark.sql import SparkSession spark SparkSession.builder.appName(check_aqe).master(local[2]).getOrCreate() print(spark.version) print(spark.conf.get(spark.sql.adaptive.enabled)) spark.stop()逻辑说明第一行spark-submit --version打印的是 Spark 版本和配套的 Scala、Java 版本用来确认环境是不是 3.x。第二段代码里的spark.conf.get是运行时读取参数比去配置文件里找靠谱。参数说明spark.sql.adaptive.enabled是 AQE 总开关在不同发行版里的默认值不一样有的默认 true有的默认 false所以不要凭记忆猜要用这行代码查。查完如果输出false就按后面第 5 章里的方式显式打开。注意AQE 不是银弹它生效有前提。比如你手动指定了repartition(1000)AQE 的动态合并分区就不会再动这个 shuffle 分区数。所以入门阶段正确的姿势是先少改参数多跑几遍看 Spark UI 里的统计再决定要不要干预。3. 本地搭一个能跑的 Spark 3.0tar 包装法、第一个提交命令学习包里的代码如果跑不起来读十遍笔记也没用。这一章给你一套能跑通本地作业的完整步骤从装环境到提交第一个 PySpark 作业每一步都是可以照着敲的。3.1 三种环境怎么选命令行、Docker、云主机先别急着下决定拿张表对比一下。方式适合场景优点缺点官方 tar 包本机安装个人入门、学习阶段环境最透明报错好排查spark-submit 直接可跑本机需要装好 JDK升级要手动换Docker 容器想隔离环境、怕搞乱本机环境干净删了重建方便镜像不一定带 Python跑 pyspark 要额外处理云主机/虚拟机多人协作、需要稳定环境配置统一资源可扩展要花钱网络和存储都要额外考虑个人入门我一般推荐第一种因为学习阶段你需要看得见 JDK、Python、Spark 各自在哪环境出问题能一步步查到原因。Docker 是好东西但等你能把本地作业跑通了再上容器也不迟。3.2 十分钟装好单机版解压、环境变量、验证版本到 Spark 官网下载页面选一个 3.x 的预编译包文件名类似spark-3.x.x-bin-hadoop3.tgz下载后按下面步骤装。前提是先确认 Java 已经装好。# 确认 Java 版本Spark 3.x 需要 JDK 8 或 11 java -version # 解压到 /opt 目录包名替换成你实际下载的文件 tar -xzf spark-3.x.x-bin-hadoop3.tgz -C /opt/ # 建一个软链接方便以后换版本 ln -s /opt/spark-3.x.x-bin-hadoop3 /opt/spark # 把环境变量写进 ~/.bashrc echo export SPARK_HOME/opt/spark ~/.bashrc echo export PATH$SPARK_HOME/bin:$PATH ~/.bashrc source ~/.bashrc # 验证是否安装成功 spark-submit --version逻辑说明-C /opt/表示把包解压到 /opt 目录下软链接的作用是以后升级 Spark 时只改链接指向不用改环境变量。source ~/.bashrc让配置立即生效。参数说明如果你的 Java 版本是 17部分 Spark 3.0 小版本会有兼容问题建议直接用 8 或 11省得后面遇到莫名其妙的反射报错。注意如果spark-submit --version报错找不到 Java先检查JAVA_HOME有没有设对很多环境问题卡在这一步。3.3 提交第一个 PySpark 作业读 CSV 做一次聚合环境装好后建一个项目目录先创建一个测试数据文件data.csv内容如下。city,age,amount beijing,25,100 shanghai,30,200 beijing,40,300 guangzhou,20,150 shanghai,35,250然后写第一个作业first_etl_demo.py功能是从 CSV 读数据过滤年龄异常值按城市统计订单金额。from pyspark.sql import SparkSession from pyspark.sql.functions import sum, col spark SparkSession.builder \ .appName(first_etl_demo) \ .master(local[2]) \ .config(spark.sql.shuffle.partitions, 4) \ .getOrCreate() df spark.read.option(header, True).csv(data.csv) # 查看列名和类型确认读进来的是什么 df.printSchema() df.show() # 过滤异常年龄age 列读进来是字符串先转整型 clean_df df.filter(col(age).cast(int) 0) # 按城市聚合amount 也做一次类型转换 result clean_df.groupBy(city) \ .agg(sum(col(amount).cast(decimal(10,2))).alias(total_amount)) \ .orderBy(total_amount, ascendingFalse) result.show() spark.stop()提交命令cd /path/to/your/project spark-submit --master local[2] first_etl_demo.py逻辑说明spark.read.option(header, True).csv(data.csv)会把第一行当表头所有字段默认读成字符串所以后面聚合前必须cast。groupBy(city)会触发一次 shufflespark.sql.shuffle.partitions控制这次 shuffle 生成多少个分区这里设成 4 是为了小数据集下看到的效果更明显。跑完后你会在终端看到按城市排序的输出这就是第一个跑通的 Spark 作业。spark-submit --master local[2]里的local[2]表示本地模式用 2 个线程模拟并发。等你把这段代码改明白再去看学习包里 day1 和 day2 的笔记会轻松很多。4. 把 1-8 天资料拆成四段学代码和笔记的正确用法拿到资料包别按文件夹顺序一天一天硬啃。Spark 学习有个规律前三天是 API 熟悉期中间两天是读写数据后三天才真正和性能与业务沾边。按下面这个拆法效率会高很多。4.1 资料包里 day1 到 day8 的常见分法天数主题动手内容完成标志day1环境搭建跑通 spark-shell 和 pysparkspark-submit --version正常输出day2RDD 与算子做一次词频统计能独立写出 map/reduceByKey 链day3DataFrame 与 SQL对 CSV 做分组聚合能用 DataFrame 和 SQL 各写一遍day4文件读写读写 CSV、JSON、Parquet知道每种格式的适用场景day5结构化流从本地 socket 读数据流跑通窗口统计并理解处理时间day6调优基础看 Spark UI、调分区数能说出某个 Job 的 shuffle 大小day7综合项目日志分析或网约车数据清洗独立完成一条清洗到聚合的流程day8复习与面试题重做前面作业加 SQL 题不看笔记能写完整代码这是最常见的内容编排你手里的包可能有出入但主线逃不开这些。注意 day5 的结构化流不是入门难点跑通一个 demo 即可别一开始就死磕 exactly-once 这类话题。4.2 笔记的读法先写“我这次要跑出什么”再看代码大多数人读笔记是从第一行读到最后一个注释合上笔记就忘。换个顺序每天开始前先在笔记空白处写三行字——今天的输入数据是什么、要得到什么结果、代码里可能有哪几个关键步骤。然后带着这三个问题去读代码。比如 day3 的笔记讲 DataFrame 聚合你先写“输入是 user_log.csv要得到每个城市每个渠道的 PV 和 UV”再去看笔记里的代码你会主动去找groupBy、countDistinct这些算子而不是被动扫一眼。读完代码后把笔记里的数据换成你自己的一个小文件改路径、改列名跑一遍。这招看似简单但它逼着你在动手之前思考把“看别人的代码”变成“验证自己的理解”。4.3 同一个需求用 RDD、DataFrame、SQL 各写一遍学习包里的代码通常只给一种写法。我建议你每个统计需求强制自己写三遍这是吃透 API 最快的方式。下面以词频统计为例。from pyspark.sql import SparkSession spark SparkSession.builder.appName(three_api_demo).master(local[2]).getOrCreate() # 准备输入两行文本 df_lines spark.createDataFrame([hello world, hello spark], string).toDF(line) df_lines.createOrReplaceTempView(lines) # 写法一RDD手动拼算子 rdd_result df_lines.rdd \ .flatMap(lambda row: row[line].split( )) \ .map(lambda word: (word, 1)) \ .reduceByKey(lambda a, b: a b) print(RDD:, rdd_result.collectAsMap()) # 写法二DataFrame API df_lines.selectExpr(explode(split(line, )) AS word) \ .groupBy(word).count() \ .orderBy(count, ascendingFalse) \ .show() # 写法三Spark SQL spark.sql( SELECT word, COUNT(*) AS cnt FROM ( SELECT explode(split(line, )) AS word FROM lines ) GROUP BY word ORDER BY cnt DESC ).show() spark.stop()逻辑说明三种写法结果一致但表达层次不同。RDD 写法的flatMap和reduceByKey是手动算子适合理解执行过程。DataFrame 写法里的explode(split(...))是 SQL 风格的列内展开引擎可以优化。Spark SQL 最接近数仓开发习惯后面做复杂指标时基本都用它。参数说明createOrReplaceTempView是注册一个临时视图作用域只在当前 SparkSession 内不落盘也不影响其他作业。如果你能把三个版本都独立写出来并且说出每个版本触发的 shuffle 位置学习包前 4 天的内容就掌握得差不多了。5. Spark 3.0 自学的五个典型坑现象、原因、解决办法这五个坑是我见过最多人踩的每个都按“现象、原因、解决”讲清楚。碰上了对照着查能少走很多弯路。5.1 pyspark 报 Python not foundexecutor 起不来现象本地跑spark-submit first_etl_demo.py日志里出现 Python worker failed to connect 或者 Python not found有时还会说找不到python3。原因Spark 的 Python worker 在 Executor 端需要找到解释器它默认找python但你机器上可能只有python3或者路径没有写进环境变量。解决显式指定解释器路径再提交。export PYSPARK_PYTHON$(which python3) export PYSPARK_DRIVER_PYTHON$(which python3) spark-submit --master local[2] first_etl_demo.pyPYSPARK_PYTHON是 Executor 端用的解释器PYSPARK_DRIVER_PYTHON是 Driver 端用的两条都设能避免本地模式下的不一致。这个坑在 Windows 上尤其常见因为python命令可能指向 Microsoft Store 的占位程序。5.2 collect() 把结果全拉回 Driver还没看数据内存就爆了现象跑一个清洗作业数据量大概几千万行最后一行代码df.collect()然后 Driver 直接OutOfMemoryError作业卡死。原因collect()会把所有分区计算结果全部拉回 Driver 端内存数据量一大必然爆。入门阶段很容易把collect()当成“看一眼结果”的默认动作。解决观察少量数据用show()和take()真的要全量遍历就先聚合再 collect。# 错误示范几千万行全部拉回 Driver rows df.collect() # 正确示范只看前 10 行观察数据结构 df.show(10, truncateFalse) df.take(10) # 需要全量统计时先聚合再 collect df.groupBy(city).count().collect()show(10, truncateFalse)里的truncateFalse表示不截断字段内容调试时能看到完整值。记住一个原则凡是执行完会返回 Driver 的操作collect、toLocalIterator都要先想想数据量。5.3 本地能读文件换 cluster 模式就报找不到路径现象本地跑spark-submit读file:///home/user/data.csv正常换成--master yarn或 Kubernetes 提交后报FileNotFoundException或者Path does not exist。原因yarn cluster模式下Driver 和 Executor 都跑在远端节点上远端节点上没有你本地那个文件路径。file://协议指向的是每个进程所在机器的本地磁盘。解决数据放到集群都能访问的存储上比如 HDFS、S3、对象存储然后修改代码里的路径前缀。如果只是一个小文件可以用--files把文件分发到 Executor 的工作目录。spark-submit --master yarn --files ./data.csv job.py代码里直接用相对文件名data.csv就能读到。注意--files分发的文件对 Driver 和 Executor 都可用不需要写完整路径但文件变了要重新提交。5.4 输出目录里全是不到 1MB 的小文件现象跑完一个分组统计写 Parquet 或 CSV 到输出目录发现生成了几百个part-xxxx文件每个才几百 KB下游加载特别慢。原因默认spark.sql.shuffle.partitions是 200每个 shuffle 分区会各自写一个文件数据量不大时自然每份都很小。大量小文件是数据湖和数仓的大忌。解决写文件前先控制分区数。# 重分区到 4 个然后写 Parquet df.repartition(4).write.mode(overwrite).parquet(output/) # 数据量特别小时合并成 1 个文件 df.coalesce(1).write.csv(output/)repartition会触发一次 shuffle把数据均匀分到 4 个分区coalesce(1)尽量不 shuffle但如果数据原本分布在几百个分区它会往一个节点凑数据反而拖慢速度。写入前先看数据量再选方案。5.5 AQE 配了却没生效参数开了像没开一样现象照着文章设置spark.sql.adaptive.enabledtrue跑完发现作业没有变快Spark UI 里的执行计划也没有变化感觉自己配了假参数。原因第一不同版本的默认值不同3.0、3.1 里 AQE 的子开关标着 experimental部分配置项要一起打开才生效第二AQE 动态合并分区对已经手动repartition的场景不生效。解决先查实际生效值再显式打开整套开关。spark-shell --master local[2] scala spark.conf.get(spark.sql.adaptive.enabled)spark-submit \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --conf spark.sql.adaptive.skewJoin.enabledtrue \ job.pycoalescePartitions.enabled负责动态合并小分区skewJoin.enabled负责处理倾斜 join这两个子开关在部分 3.0 小版本里要手动配合总开关一起打开。判断是否生效去 Spark UI 的 SQL 标签页看执行计划里有没有AdaptiveSparkPlan节点有才是真的开了。6. 学完 1-8 天之后用三个办法验证自己是否真的入门学习包看完不代表掌握下面三个验证方法都不花钱但能真实反映你能不能上手干活。6.1 找一份真实 CSV完整跑一遍清洗到聚合去网上找一份真实的业务数据网约车订单、电商日志都行只要字段够多、有脏数据。要求自己独立完成读入、过滤异常值、维度聚合、写出 Parquet并且 Spark SQL 和 DataFrame API 各写一遍。中间不看任何笔记。能独立跑通说明 day2 到 day4 的内容是真学会了如果卡在某个算子或类型转换上恭喜你那里就是你的薄弱点。6.2 学会从 Spark UI 读结论而不是只看“跑通”跑完作业后打开http://localhost:4040这里不会骗你。重点看 Stage 页面和 Executor 页面。现象可能原因某个 Stage 的 Shuffle Read 超过 1GB分区数过多或 join key 选择不佳单个 Task 耗时远超中位数数据倾斜考虑 AQE 或加盐Executor 的 GC 时间占比高内存不足加大 executor 内存或减少缓存一句话跑通只是开始能解释为什么跑得快才是入门。6.3 给自己建一张参数速查卡把学习过程中用过的参数记到一张固定表格里以后照着查。参数作用常用值spark.sql.shuffle.partitions控制 shuffle 分区数集群 CPU 核数的 2-3 倍spark.sql.adaptive.enabledAQE 总开关truespark.executor.memory单个 executor 内存根据集群资源定别超节点内存spark.executor.cores单个 executor 可用核数2-4spark.sql.autoBroadcastJoinThreshold小表自动广播的阈值默认 10MB视内存调整我一般会在本地维护三个固定文件env.sh放环境变量submit.sh放提交参数notes.md记踩坑记录。换新环境先跑一次spark-submit --version提交作业只改submit.sh里的路径遇到问题先在notes.md里翻旧账。这套习惯让我从第 5 天就想放弃的状态变成真正敢说自己入门了。希望帮到你。本文还有配套的精品资源点击获取