ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spark 3.0入门到调优:RDD、DataFrame与Structured Streaming实战笔记

Spark 3.0入门到调优:RDD、DataFrame与Structured Streaming实战笔记 简介围绕Spark 3.0.1的入门到精通学习资源主要面向大数据初学者和希望系统进阶的开发者以2021贺岁课程为蓝本覆盖环境搭建、Spark Core、Spark Streaming、Spark SQL、Structured Streaming、综合案例、多语言开发、3.0新特性与性能调优等九大主题构成一条从基础概念到实战调优的完整学习路径。资源包共244个文件总大小86.9MB以png截图、md笔记、scala源码为主并附带9个zip子包和1个json配置。其中217张截图直观展示安装与运行界面10份markdown笔记归纳每节重点7个scala源码文件供动手实践子zip包则可进一步拆分课程素材便于按需学习。已有609人学习参考。按1至8天组织每天配套独立笔记与代码可配合课程搭建Spark集群理解RDD与DataFrame等核心抽象掌握流式处理与结构化查询的开发方式同时了解3.0版本的自适应查询执行、动态分区裁剪等新特性及常用调优方法。对于零基础入门或准备系统梳理Spark知识的学习者这套资源既提供阶段化指引又能通过代码与截图的结合有效降低上手门槛适合自学及复习巩固。1. Spark 3.0 大数据入门这份 8 天笔记为什么值得照着敲在聊大数据学习路线的时候很多人第一站就卡在 Spark 上Hadoop 那套还没吃透Spark 又抛来 RDD、DataFrame、Streaming 一堆概念真到集群上跑又全是环境问题。这份课程用的不是老版本而是官方 2020 年 9 月 8 日发布的 Spark 3.0.1 稳定版从 Spark 环境搭建一路打到性能调优8 天的代码和笔记全在一起算是一条能把大数据入门到精通串完整的路径。适合刚学完 Hadoop、准备啃 Spark 的初学者也适合想快速接手 Spark 项目但缺一份完整代码参考的开发。整套资源里既有按天拆好的 Spark-day01 到 Spark-day06 笔记也有独立的 Spark 综合案例、Spark 多语言开发以及配套的 question_info.json照着敲比自己瞎试快得多。2. Spark 环境搭建与 SparkCore先跑通第一个分布式任务再谈优化2.1 版本选型为什么锁 Spark 3.0.1 而不是旧版拿到这份资源第一步不是急着写代码而是把 Spark 3.0.1 的运行环境对齐。这个版本是 Apache 官方 2020 年 9 月发布的 3.0 系稳定版本编译基础是 Scala 2.12官方对 Hadoop 的兼容范围覆盖 2.7 和 3.2。JDK 方面Spark 3.0 要求 8 或 11我自己习惯选 JDK 8u202 以上版本因为后续 Spark SQL 和 Hive 元数据交互时旧版 JDK 偶尔会抛奇怪的反射异常排查起来非常浪费时间。# 我习惯的目录规划多套版本切换时不互相污染 mkdir -p /opt/bigdata/spark /opt/bigdata/hadoop /opt/bigdata/jdk tar -zxvf spark-3.0.1-bin-hadoop2.7.tgz -C /opt/bigdata/spark mv /opt/bigdata/spark/spark-3.0.1-bin-hadoop2.7 /opt/bigdata/spark/spark-3.0.1这段命令的作用是解压 Spark 安装包并把长目录名改成短路径。这里有个选型细节容易被忽略下载时一定选名字里带bin-hadoop2.7或bin-hadoop3.2字样的预编译包别下src结尾的源码包。源码包拿回来要用 Maven 重新编译时长半小时起步而且编译环境不一致时spark-shell直接起不来属于最典型的入门回收站行为。环境变量是另一个高频翻车点。Spark 3.0.1 启动时同时读JAVA_HOME和HADOOP_HOME少配一个就报找不到可执行文件。我一般把配置写进/etc/profile.d/spark.sh里这样切用户之后还能正常识别。export JAVA_HOME/opt/bigdata/jdk/jdk1.8.0_202 export HADOOP_HOME/opt/bigdata/hadoop/hadoop-3.2.1 export SPARK_HOME/opt/bigdata/spark/spark-3.0.1 export PATH$PATH:$JAVA_HOME/bin:$HADOOP_HOME/bin:$SPARK_HOME/bin配置完执行spark-shell --version能打出版本号环境就算通了一半。接下来马上进入运行模式选型这块在 Spark-day01.md 里有专门对照表我这里也列了一份Master 参数适用阶段注意事项local[*]课程前三天的单机练习数据量小于 2GB 时够用最大并发受本机核数限制spark://host:7077standalone 模拟集群需要先启动 master 和 worker适合学习集群调度机制yarn生产环境依赖 HDFS适合和 Hadoop 生态一起用前三天用local[*]完全够没必要一上来就搭三台虚拟机。我见过不少人第一天就在配集群结果环境没弄好WordCount 都没跑通信心先被打了一半。2.2 RDD 核心算子从 WordCount 理解完整数据流转Spark-day01 和 Spark-day02 里WordCount 是反复出现的例子。它看起来简单却把 RDD 的三个关键特性全带出来了不落盘、懒加载、分区并行。理解这三件事后面调优时才有抓手。from pyspark import SparkContext, SparkConf conf SparkConf().setAppName(wordcount_demo).setMaster(local[*]) sc SparkContext(confconf) lines sc.textFile(file:///data/input/words.txt, minPartitions2) words lines.flatMap(lambda line: line.split( )) pairs words.map(lambda word: (word, 1)) counts pairs.reduceByKey(lambda a, b: a b) counts.saveAsTextFile(file:///data/output/wc_result)代码逻辑分四层textFile读取文本并按参数切分分区flatMap把每行打散成单词map把单词映射成(word, 1)二元组reduceByKey先做分区内局部聚合再做跨分区全局聚合。saveAsTextFile有个强制约束输出目录必须不存在否则直接跑FileAlreadyExistsException这是 Spark 对输出路径的保护机制踩过一次就不会忘。minPartitions2值得单独讲。它是期望的最小分区数不是绝对保证实际分区数还看文件存储方式本地文件系统默认按 32MB 切块HDFS 按 128MB 切块。从这个小参数能引出调优课里的一个重要原则分区数太少单个任务处理数据量过重分区数太多调度开销反而拉低吞吐。常见经验值是让每个分区承载约 128MB 数据量。2.3 DataFrame 转换Catalyst 优化器带来的思维切换到 Spark-day02 后半部分笔记开始引入 DataFrame。这里要做一个思维切换RDD 关注的是数据如何被分散处理DataFrame 关注的是表结构上的声明式操作底层 Catalyst 优化器会自动把filter、join的执行顺序重排到更优路径。import org.apache.spark.sql.SparkSession val spark SparkSession.builder() .appName(rdd_to_df) .master(local[*]) .getOrCreate() import spark.implicits._ val rdd spark.sparkContext.textFile(file:///data/input/sales.txt) val df rdd.map { line val fields line.split(,) (fields(0), fields(1).toDouble) }.toDF(category, amount) df.filter($amount 100) .groupBy(category) .agg(sum(amount).as(total_amount)) .show()这段代码有两个隐性知识点。第一次用toDF的人十有八九会漏掉import spark.implicits._导致toDF和$amount编译报错这是 Scala 隐式转换机制决定的。第二个坑是 DataFrame 的算子全部返回新对象filter之后原df不受影响这和 RDD 的不可变语义一致但刚接触时总有人觉得应该像 Python 列表那样原地修改。Day02 笔记里有一段 RDD 和 DataFrame 在执行计划层面的对比。RDD 的filter是直接作用于分区的函数调用DataFrame 的filter先生成逻辑计划节点Catalyst 再结合统计信息决定是否做谓词下推。这种优化在 Parquet 列式存储上效果尤其显著后面的综合案例做数据清洗时会反复用到。从这个角度看Spark Core 阶段真正要练的不是算子数量而是理解数据流转的直觉数据从哪来经过哪些转换最后落到哪。2.4 序列化与依赖本地跑通不等于集群能跑通单机用local[*]跑通 WordCount 只是第一步把任务真正丢到集群上时碰到的第一个大坑基本是序列化。Spark 3.0.1 在集群模式下driver 把任务分发给 executor 时闭包里引用的对象要通过 Java 序列化传输。自定义类不实现Serializable提交任务时就会报Task not serializable。class MyFilter extends Serializable { val threshold: Double 100.0 def isHigh(amount: Double): Boolean amount threshold } val myFilter new MyFilter() df.filter(row myFilter.isHigh(row.getAs[Double](amount))).count()解决方案是让闭包中引用的类实现Serializable或者把常量提取成方法内局部变量避免闭包捕获整个对象。Spark-day02.md 对这个问题有单独标注因为综合案例里自定义 UDF 时还会再踩一次。记住一个判断标准凡是 driver 端定义、executor 端使用的对象都必须过序列化这一关。3. Spark SQL 与流处理从 DataFrame API 到 Structured Streaming 的实战配置3.1 Spark SQL 读写通道Parquet、JDBC 与 Hive 表的打通Spark SQL 这部分课程内容核心不是 SQL 语法本身而是读写通道怎么打通。Spark 3.0.1 的 SQL 可以读 HDFS 上的 Parquet、JSON、CSV也能通过 JDBC 连 MySQL还能挂上 Hive 元数据服务。实操里我优先选 Parquet原因是列式存储加压缩配合 Catalyst 的谓词下推全表扫描的量能砍掉一大截。# PySpark 写法读取 Parquet 后注册临时表再跑纯 SQL from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(spark_sql_demo) \ .config(spark.sql.adaptive.enabled, true) \ .getOrCreate() df spark.read.parquet(hdfs:///data/input/sales.parquet) df.createOrReplaceTempView(sales) result spark.sql( SELECT category, count(*) AS order_cnt, round(sum(amount), 2) AS total_amount FROM sales WHERE dt 2024-01-15 GROUP BY category ORDER BY total_amount DESC ) result.show()这个例子里createOrReplaceTempView是关键动作它把 DataFrame 注册成临时表之后就能用纯 SQL 操作。Spark 3.0 里临时表默认只在当前 SparkSession 内可见退出就消失想做跨会话共享得用saveAsTable落成持久表。CSV 读数的坑也值得单独说。CSV 没有 Schema 信息Spark 默认会把所有列当字符串读后面做聚合时才发现类型不对。# 指定 schema 避免类型推错 from pyspark.sql.types import StructType, StructField, StringType, DoubleType schema StructType([ StructField(category, StringType(), True), StructField(amount, DoubleType(), True), StructField(dt, StringType(), True) ]) df spark.read \ .option(header, true) \ .schema(schema) \ .csv(hdfs:///data/input/sales.csv)这里header选项控制首行是否作为列名schema显式声明列类型。不写 schema 临时能跑但数据量一大或者源文件里混入空值类型推错的问题会延迟到聚合阶段才爆出来那时候排查成本就高了。日志和数据质量验证应该在数据接入的那一刻就做而不是等结果对不上再去翻数据。3.2 Spark Streaming 与 Structured Streaming两代流处理的选型Spark-day03 和 day04 开始进入流处理这里要分清两个概念Spark Streaming 是早期的 DStream 模型按微批处理API 和 RDD 一脉相承Structured Streaming 是 Spark 2.0 之后引入的基于 DataFrame 的流处理模型到 3.0 已经进入成熟期支持mapGroupsWithState、flatMapGroupsWithState这些有状态算子。选型上新项目直接学 Structured Streaming除非要维护老代码。理由很现实Structured Streaming 的 API 和 Spark SQL 几乎一致能复用批处理里的思维而且 3.0 的 Streaming 和 SQL 的执行引擎已经统一到同一套 Catalyst 优化路径上。# Structured Streaming 读取 Kafka 的典型写法 from pyspark.sql import SparkSession from pyspark.sql.functions import window, col spark SparkSession.builder \ .appName(streaming_kafka) \ .config(spark.sql.shuffle.partitions, 8) \ .getOrCreate() df spark.readStream \ .format(kafka) \ .option(kafka.bootstrap.servers, localhost:9092) \ .option(subscribe, user_actions) \ .option(startingOffsets, latest) \ .load() parsed df.selectExpr(CAST(value AS STRING) as json) \ .selectExpr(from_json(json, user_id STRING, action STRING, ts TIMESTAMP) as data) \ .select(data.*) windowed parsed \ .withWatermark(ts, 10 minutes) \ .groupBy(window(ts, 5 minutes), action) \ .count() query windowed.writeStream \ .outputMode(append) \ .format(console) \ .trigger(processingTime5 seconds) \ .start() query.awaitTermination()这段代码里值得展开的点withWatermark设置水位线用来容忍乱序数据window(ts, 5 minutes)定义滚动窗口trigger(processingTime5 seconds)控制微批间隔。startingOffsets设为latest表示只消费启动后的增量数据调试阶段这样最快生产环境一般用earliest配合 checkpoint 做断点续跑。3.3 流式任务的 checkpoint 与语义至少一次与精确一次的差别流处理项目里checkpoint 不是可选配置是必选。Structured Streaming 把 Offset 和状态数据都写进 checkpoint 目录用来做故障恢复。不配 checkpoint任务只要重启offset 就找不回来了。query windowed.writeStream \ .format(parquet) \ .option(path, hdfs:///data/output/stream_result) \ .option(checkpointLocation, hdfs:///data/checkpoint/stream_result) \ .outputMode(append) \ .trigger(processingTime10 seconds) \ .start()checkpointLocation必须绑定一个不和其他任务共用的目录。如果两个流任务共用一个 checkpoint 路径重启时会因为目录里已经有 driver 元数据而直接拒绝启动。这个坑我在联调时遇到过现象是日志一直刷A checkpoint directory must not be used by multiple applications原因就是复制粘贴时没改路径。处理语义方面Spark 3.0.1 的 Structured Streaming 默认提供至少一次语义要精确一次需要额外配合 Kafka 的幂等写出和事务日志。对入门阶段来说先保证至少一次再把重复数据在聚合层做去重比一上来就追求精确一次更实际。4. Spark 性能调优与常见问题排查从执行计划到 5 个真实现场4.1 用 Web UI 和物理执行计划定位瓶颈课程后半部分反复强调调优之前先看 UI。Spark 3.0.1 的 Web UI 默认跑在 4040 端口包含 Jobs、Stages、Storage、Executors 四个关键页签。看到一个耗时异常长的 Stage点进去看每个 Task 的处理时间分布基本上能判断是数据倾斜还是资源不足。-- 查看物理执行计划SQL 里加 EXPLAIN EXPLAIN COST SELECT category, sum(amount) FROM sales WHERE dt 2024-01-15 GROUP BY category;EXPLAIN COST会输出优化后的逻辑计划和带行数估算的物理计划。如果发现Exchange节点两侧数据量差距过大基本可以确定是 join 键分布不均导致的倾斜。另一个常用手段是把计划打到日志里用df.explain(true)在 DataFrame API 里直接看完整计划。4.2 资源参数配置executor 数量、内存与并行度怎么定Spark-day07 性能调优笔记里最实用的内容是一套资源参数的初始值设定方法。很多新手上来就调spark.executor.memory但忽略了 executor 数量和并行度之间的配合导致内存够用但 CPU 利用率极低。# 常见 submit 参数模板 spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 20 \ --executor-cores 4 \ --executor-memory 8g \ --driver-memory 4g \ --conf spark.sql.shuffle.partitions200 \ --class com.example.MainJob /data/app/spark-demo.jar关键参数的含义我按经验列一张表参数作用调整依据--num-executorsexecutor 总数根据队列资源和数据量定不是越多越好--executor-cores每个 executor 的 CPU 核数建议 2-5核数过多会让 HDFS IO 成瓶颈--executor-memory单 executor 可用内存留出 20%-30% 做 JVM 开销和存储spark.sql.shuffle.partitionsshuffle 后的默认分区数数据量大就调大通常 200-500一个容易踩的坑是 executor 内存只算堆内没算堆外。spark.executor.memoryOverhead需要单独留尤其在跑 Python UDF 和加解密类任务时堆外内存不足会抛Container killed by YARN for exceeding memory limits。出现这个错误优先加memoryOverhead而不是继续加堆内内存。4.3 五个常见问题的现象、原因、解决问题一java.lang.OutOfMemoryError: GC overhead limit exceeded现象任务跑到 shuffle 阶段频繁 Full GC然后直接 OOM。原因executor 内存中spark.memory.storage和spark.memory.execution存在动态抢占默认各占 50%但如果 RDD cache 数据过量又同时跑大聚合内存就被挤爆了。解决先检查代码里是否有过量的cache()或persist()及时unpersist()再调spark.memory.fraction从默认 0.6 适当降到 0.5给 JVM 留更多余量最后把 executor 内存加大一到两档。问题二java.lang.IllegalArgumentException: Size exceeds Integer.MAX_VALUE现象groupByKey或reduceByKey在单个分区内累积的数据超过 2GB。原因数据倾斜导致某个 key 的数据量巨大单分区无法容纳。解决把这个算子改成两阶段聚合先加随机盐键打散再聚合去盐。课程综合案例里有个按用户维度聚合的场景就是这么处理的。问题三Caused by: java.io.FileNotFoundException: File does not exist现象任务重启后读 HDFS 上的临时文件失败。原因Spark 在 shuffle 过程中生成的临时文件被定期清理任务重试时旧文件句柄已经失效。这个问题在开启动态资源分配时尤其常见。解决关闭动态资源分配spark.dynamicAllocation.enabledfalse或者把spark.shuffle.service.enabledtrue开启让 shuffle 文件由外部服务托管。问题四org.apache.spark.SparkException: Task not serializable现象提交任务时直接报错栈指向自定义过滤类。原因闭包捕获了不可序列化的对象。我在 2.4 节踩过一模一样的坑。解决类实现Serializable或把成员变量提取成局部变量传入函数。问题五argument type mismatch或Column xxx does not exist现象DataFrame 的列名和查询不一致运行时报错。原因CSV 或 JSON 数据源列名大小写不一致或者读入时列名被 Spark 改了比如带逗号的列名会被加上反引号。解决统一在数据接入层做列名清洗用withColumnRenamed映射成规范命名再用df.printSchema()验证。这个习惯能省掉后续 SQL 阶段大量莫名其妙的报错。5. Spark 综合案例与多语言开发把 8 天内容串成一个可复现项目5.1 综合案例拆解从数据清洗到指标输出的完整链路Spark 综合案例.md 是整个课程里最值得反复读的文件。它不只是写几段代码而是把用户行为数据从原始日志到最终报表的完整处理链路串了起来。链路的每一段正好对应前面 8 天的某个知识点数据接入对应 Spark SQL实时清理对应 Structured Streaming离线统计对应 SparkCore最终输出又落回 Spark SQL 写 Hive 或 Parquet。实际去做这类案例时我一般按四步拆先定义需求即最终要输出哪些表和指标再定数据源分清离线文件、Kafka 流和数据库表然后画数据流标出每个转换节点的输入输出最后才写代码。整个链路里最容易出问题的步骤是数据清洗空值处理策略不统一会导致结果数据凭空多出很多脏行。# 综合案例中常见的数据清洗片段 from pyspark.sql.functions import col, when, isnan df_clean df \ .dropDuplicates([user_id, action, ts]) \ .filter(col(user_id).isNotNull()) \ .withColumn(amount, when(col(amount).isNull(), 0).otherwise(col(amount))) \ .withColumn(dt, col(ts).substr(0, 10))这段清洗逻辑做四件事按业务键去重、过滤空用户、把缺失金额补零、从时间戳截取日期分区。注意dropDuplicates去重依赖字段选择去重键必须能唯一标识一条记录否则会把正常的多条行为误删。5.2 多语言开发Scala、Python、Java 三套 API 的协作边界Spark-多语言开发.md 这份笔记解决了一个很现实的问题团队里有人用 Python有人用 Scala最终任务怎么统一管理。Spark 3.0.1 本身是 JVM 应用Scala 是原生 API性能和功能覆盖最全Python 通过 PySpark 调用底层走 Py4J 桥接Java API 在 3.0 里已经不那么常用但维护老项目时还会遇到。下面这份表格是我从笔记里归纳的选择标准语言适用场景注意点Scala核心 ETL、性能敏感任务、UDF编译期能发现大部分类型问题Python数据探索、机器学习模型调用、快速原型UDF 序列化开销大尽量用内置函数Java存量系统集成、特定团队技术栈API 比 Scala 繁琐开发效率偏低用 PySpark 时最需要注意的是 Python UDF 的性能问题。每调用一次 Python UDF数据都要在 JVM 和 Python 进程之间做序列化传输。能用pyspark.sql.functions内置函数解决的就不要自定义 UDF性能差距能达到数倍。Spark 3.0 的 Pandas UDF 可以缓解一部分性能问题但要额外处理依赖包分发入门阶段可以先了解不急着深入。5.3 从案例代码反推真实项目分层综合案例看完了接下来要把代码结构映射到真实项目的分层上。我见过太多人学完 Spark 之后所有逻辑都堆在一个main方法里能跑但是改不动。真实项目通常分四层接入层负责读各种源数据统一格式和 Schema清洗层空值处理、去重、格式规范化计算层指标逻辑、聚合规则输出层写 Hive 表、Parquet 文件或 Kafka 消息// 分层写法的骨架示意 object UserActionJob { def main(args: Array[String]): Unit { val spark SparkSession.builder().appName(user_action_job).getOrCreate() val raw DataLoader.loadUserActions(spark, args(0)) val cleaned DataCleaner.clean(raw) val metrics MetricsCalculator.computeDailyMetrics(cleaned) MetricsWriter.writeToHive(metrics, dws.user_action_daily) spark.stop() } }每个模块一个单例对象职责单一后面替换数据源或者改清洗规则时只动对应模块就行。这个习惯比任何调优参数都重要可维护的代码结构才是大数据项目能持续迭代的基础。6. 进阶验证Spark 3.0 新特性和调优后的自查习惯6.1 用 AQE 和动态分区修剪验证 3.0 新特性Spark 3.0 相比 2.4 最明显的进步是自适应查询执行简称 AQE。以前spark.sql.shuffle.partitions设成 200所有 shuffle 都按 200 个分区跑不管数据量多小。开了 AQE 之后Spark 会在运行时根据实际 shuffle 数据量自动合并小分区这个特性在 3.0 里默认开启我一般会主动验证一下效果。spark-submit --conf spark.sql.adaptive.enabledtrue --conf spark.sql.adaptive.coalescePartitions.enabledtrue --class com.example.MainJob /data/app/spark-demo.jar验证方法很简单跑完任务去 Web UI 的 Stages 页面看第二个 Stage 的 Task 数量如果明显少于spark.sql.shuffle.partitions的设定值就说明 AQE 生效了。另外尝试把spark.sql.adaptive.coalescePartitions.initialPartitionNum设成一个预期值观察执行计划的变化能加深对优化机制的理解。6.2 任务上线前固定的五个检查动作从那以后我每次提交 Spark 任务都强制自己走一遍这套检查先看输入路径是否存在、Schema 是否匹配再看 SQL 执行计划里有没有全表扫描和过度 Shuffle然后确认 checkpoint 目录没被其他任务占用接着检查输出目录有没有历史残留避免FileAlreadyExistsException最后跑一个count()验证数据量符合预期。这套习惯救过我不少次。最典型的一次是凌晨上线定时任务忘了清理前一天留下的输出目录结果任务一开始就失败等到早上才发现。现在这些检查全写进启动脚本里谁跑都一样。把这套验证方法连同前面的调优参数、踩坑记录结合起来再回去看 Spark-day01 到 Spark-day06 的笔记会发现每一章的细节都串起来了。资源不难难的是把坑提前踩完。希望帮到你。本文还有配套的精品资源点击获取
返回列表