ARTICLE DETAIL

资讯详情

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

Spark 3.0入门实战:RDD到DataFrame的8天学习路线与调优

Spark 3.0入门实战:RDD到DataFrame的8天学习路线与调优 简介面向Spark入门学习者的大数据实战资料包覆盖Spark 3.0.1从环境搭建到性能调优的完整学习路径。资源围绕Spark Core、Spark Streaming、Spark SQL、Structured Streaming、多语言开发及3.0新特性等核心模块整理了1至8天的课程代码与笔记适合希望系统掌握Spark框架并应用于实际项目的初学者或数据工程师。压缩包共244个文件大小约86.9MB以217张原理图解、10篇Markdown笔记、9个示例代码压缩包、7个Scala源文件及1个JSON配置文件为主要构成既可用于对照学习也可直接参考项目结构。目前已有609人学习下载。资料不仅提供分日期的代码与Markdown解析还包含综合案例与性能调优专项内容能帮助读者减少搜索成本按章节循序渐进地构建Spark知识体系。1. 拿到 Spark 3.0 入门包里最先该做什么8 天学习路线的真实节奏拿到“大数据入门 spark3.0 入门到精通 1-8day 代码-笔记.zip”这种资源的人九成是把它丢进网盘吃灰。我见过太多人打开压缩包看到 Day1 到 Day8 的目录扫一眼就关掉因为不知道从哪下手。其实这 8 天的划分本身就是一条能落地的大数据学习路线前两天是核心抽象和运行模式中间两天是环境与代码后四天是 Spark SQL 清洗、分析和调优。这个包解决的是零基础到能独立提交一个 Spark 任务的全过程适合大数据专业学生拿来做毕设选题前的技术摸底也适合刚转大数据开发、想看清楚 Day 后面的代码和笔记怎么串起来的人。真正会用这套东西的人不是翻笔记而是按 Day 顺序把代码跑一遍、再对照笔记改参数。2. 从 RDD 到 DataFrameSpark 3.0 里必须分清的三个抽象2.1 RDD 为什么不是首选从 3.0 的优化器倒推Day1 的笔记几乎都会先讲 RDD但很多初学者学完就懵了明明最后写业务代码都用 DataFrame为什么还要花一整天学一个“快被淘汰”的东西我的看法是RDD 不是要被淘汰而是它的定位变了。在 Spark 3.0 里RDD 是底层执行模型DataFrame 才是面向业务的上层接口但你不理解 RDD 的分区、依赖和 Shuffle后面遇到数据倾斜、Executor 挂掉这类问题连排查方向都没有。顺着大数据架构的四个层次去看存储层、计算引擎层、资源调度层、应用层Spark 属于计算引擎层而 RDD 就是引擎内部的数据结构。Spark 3.0 之所以推荐 DataFrame是因为它带着 SchemaCatalyst 优化器能对执行计划做谓词下推、列剪枝、常量折叠Tungsten 还能把数据编码成二进制格式走堆外内存和代码生成。RDD 没有 Schema优化器无从下手只能按用户写的一步步执行。所以同一个文件用 DataFrame 读比用 RDD 读通常快 30% 到 50%数据越宽越明显。学习路线里 Day1 讲 RDD 的真正目的是让你理解 Spark 的骨架。比如你写了一个rdd.groupByKey()它背后会发生什么——按 key 哈希、分区、落盘、拉取、聚合这一整套叫 Shuffle。你理解了 Shuffle才能理解为什么别人说“尽量用reduceByKey而不是groupByKey”reduceByKey在 map 端先做一次预聚合Shuffle 的数据量小很多。笔记里如果只讲 API 不讲这层后面的 Day 全部会变成“照着敲但不知道为啥”。我自己带人入门时会让对方先看一天的 Spark 源码目录结构再回去看笔记里的 RDD 章节感受完全不一样。你不需要读懂全部源码但至少要知道org.apache.spark.rdd下面的类和算子对应哪些作业行为。等学到 Day5 或者 Day6 再回来补这个认知代价就大了。2.2 DataFrame 与 Dataset代码上的差别Day4 的笔记里同一条数据通常会给出两种写法。我拿一个最常见的 JSON 日志解析来对比你先感受下差别。// RDD 写法自己解析字段自己处理脏数据 val rdd sc.textFile(hdfs:///logs/app.json) .map { line val arr line.replace({, ).replace(}, ).split(,) (arr(0).split(:)(1).trim, arr(1).split(:)(1).trim.toInt) } .filter(_._2 0) // DataFrame 写法声明 Schema交给框架 val df spark.read.json(hdfs:///logs/app.json) .filter(count 0) .select(appName, count)这个对比在 Day4 的代码包里应该也有类似版本。RDD 写法的每一行都要自己控制去花括号、切分、转类型、过滤逻辑没错但性能和可维护性都不好。DataFrame 写法的关键是spark.read.json能自动推断 Schemafilter(count 0)是列名直接过滤连解析都省了。Catalyst 优化器看到select会把不需要的列剪掉读文件时只反序列化需要的字段这种优化在 RDD 版里完全不存在。这里有个参数值得注意spark.sql.adaptive.enabledSpark 3.0 默认开。它能在运行时根据实际 Shuffle 数据量动态调整分区数比如你设了 200 个分区但数据只有 1GB它会把分区数降下来减少每个 Task 的空转开销。这个特性在 2.4 还是实验性的3.0 变成默认开启也是我从 2.x 迁到 3.0 后体感最明显的变化之一。Dataset 是强类型版本在 Scala 里写case class配.as[]使用Java 里更常用。入门阶段可以先放着把 RDD 和 DataFrame 的关系理清就够应付学习路线前 5 天的内容了。你要记住的只是RDD 是底层实现DataFrame 是上层接口两个 API 能互相转换df.rdd随时可以拿到底层 RDD 去做自定义操作。2.3 宽依赖与窄依赖决定 Shuffle 与运行速度的底层逻辑Day2 笔记里有个必考概念叫依赖类型。窄依赖是说每个父分区最多被一个子分区使用比如map、filter、union这类操作不需要跨节点传数据在同一个 Executor 里算完就行。宽依赖是父分区被多个子分区使用比如groupByKey、join、repartition它必然触发 Shuffle也就是要把数据按 key 重新分布到多个节点上。这个知识点直接关系到大数据的集群部署策略。你在本地单机跑感受不到什么一旦上了集群宽依赖就是整个作业的瓶颈点。举个例子两张表 join如果按键分区恰好一致就能走窄依赖的优化数据不用移动如果分区策略不一致就要全量 Shuffle 一次几十 GB 的数据在节点间传一遍时间从分钟级变成小时级。我一般判断一份入门笔记写得好不好就看它有没有讲这个点。因为 Day5 之后你遇到的 90% 的性能问题最后都能归到“这里不该有宽依赖”或者“这个宽依赖的数据量太大了”。比如一个 join 之后跟着distinct再跟着groupBy中间有多少次 Shuffle能不能把distinct提前到 join 之前做用窄依赖过滤掉重复数据让 join 的输入变小这些取舍全靠对宽窄依赖的敏感度。3. 本地先跑通最小 Spark 3.0 作业环境搭建与第一条命令3.1 装什么、怎么装JDK/Scala/Spark 的版本对应Day3 的环境搭建我遇到的最多问题是“装完 Spark 3.0 启动就报错”。绝大多数情况不是 Spark 本身的问题而是它和 JDK、Scala 的版本不匹配。Spark 3.0.x 的发布包基于 Scala 2.12 编译所以本地机器的 Scala 环境最好对齐 2.12JDK 用 8 或者 11我自己一直用 JDK8因为生产集群上 JDK 8 的兼容性最稳。组件推荐版本备注JDK811 也行但排查问题资料少Scala2.12.x必须和 Spark 编译版本一致Spark3.0.x官方包自带 Scala 类库无需单独装 ScalaHadoop 客户端3.2.x 兼容本地跑local[*]可以不要这里有个关键是本地学习其实可以不装 Scala直接用 Spark 自带的spark-shell和 pyspark 就行因为它们已经把 Scala 类库打进去了。只有用 sbt 或 Maven 自己写工程打包时才需要修改pom.xml里的 scala.version 和 spark.version。我做集群之前会在本地把 Day4 的代码用spark-shell跑通一遍再进 IDE 写工程两个步骤分开排错时少一半变量。3.2 spark-shell 快速验证读本地文件的最小命令装好之后先别急着写代码用一条命令验证安装是否成功spark-shell --master local[2]--master local[2]的意思是让 Spark 直接在本地进程里跑使用 2 个 CPU 核。2 这个数字不是随便写的它能让你本地的执行行为更接近集群的多 Task 并行而不是一个核从头串到尾。进入 scala 交互式环境后会输出 Spark 版本和 Web UI 地址一般是http://localhost:4040。然后跑一个最经典的 WordCount验证数据读取和算子都正常val rdd sc.textFile(file:///data/test.log) val wordCount rdd .flatMap(_.split( )) .map(word (word, 1)) .reduceByKey(_ _) wordCount.collect.foreach(println)这段代码的逻辑一看就懂flatMap把整行文本拆成单词map给每个单词配一个计数 1reduceByKey把相同单词的计数加起来。我特意用file://前缀是因为sc.textFile(data/test.log)默认去 HDFS 里找路径没有 HDFS 环境它会一直卡在连接超时上。参数上值得改的是spark.default.parallelism本地模式默认是机器核数如果你觉得 Task 碎片化严重可以在启动时加一句--conf spark.default.parallelism4把分区数提前固定下来。3.3 打包提交到集群spark-submit 的参数怎么给Day5 的内容开始涉及集群提交。你不可能永远在 spark-shell 里写交互式代码生产上都是把工程打包成 jar 再用spark-submit提交。我用 Maven 时最简配置是pom.xml里加 spark-core 和 spark-sql 两个依赖scope 设为 provided这样打出来的包体积小提交时 Spark 会用自己的类库。提交命令我一般这样写spark-submit \ --class com.example.WordCount \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --num-executors 4 \ --executor-cores 2 \ target/wordcount-1.0.jar \ /data/input /data/output参数含义逐个说--class指定 main 方法所在类--master yarn表示提交到 YARN 资源队列--deploy-mode cluster让 driver 也跑在集群里而不是本地--executor-memory 4g给每个 Executor 分配堆内存--num-executors 4是申请 4 个 Executor--executor-cores 2是每个 Executor 用 2 个核最后的两个路径是传给 main 函数的输入输出参数。这套参数不是死板的。作业是内存密集还是 CPU 密集决定你要调内存还是调核数集群资源紧张时num-executors超过队列上限会被直接拒绝。我自己的习惯是第一天先把executor-memory和executor-cores调小跑通流程确认代码没问题再逐步加资源避免一次申请太多被运维找上门。4. Spark SQL 在入门包里的实战场景从数据清洗到数据分析4.1 网约车订单数据清洗pyspark 代码与四个检查点Day6 的里程碑在入门包里一般是一个真实场景的项目。短视频平台的教学案例里网约车订单分析出现频率最高其次是校园大数据类的数据清洗与可视化。这类项目的共通点是CSV 又脏又乱需要自己去空值、去重复、修格式、过滤异常。我用 pyspark 给你演示一段能直接跑通的数据清洗主线。from pyspark.sql import SparkSession from pyspark.sql.functions import to_timestamp, col, count spark SparkSession.builder \ .appName(order_clean) \ .master(local[2]) \ .config(spark.sql.shuffle.partitions, 8) \ .getOrCreate() df spark.read.csv( /data/order.csv, headerTrue, inferSchemaTrue ) # 四个清洗动作去重 / 过滤 / 时间规范化 / 空值处理 clean_df df \ .dropDuplicates([order_id]) \ .filter(col(amount) 0) \ .withColumn(order_time, to_timestamp(col(order_time), yyyy-MM-dd HH:mm:ss)) \ .dropna(subset[city_id, amount]) clean_df.write.mode(overwrite).parquet(/data/order_clean.parquet)这段代码的学习重点在四个清洗动作。dropDuplicates([order_id])按订单号去重如果一条订单在采集链路里被重复上报这一行就是止损点filter(col(amount) 0)把金额为负的退款记录和异常值直接过滤掉to_timestamp把字符串时间转成标准时间类型这一步不做后面的按小时聚合全是乱的dropna(subset...)只删关键字段为空的记录不删整个表。最容易忽略的配置是.config(spark.sql.shuffle.partitions, 8)。清洗阶段会有dropDuplicates这类宽依赖算子它的 Shuffle 分区数默认是 200但本地跑 200 个分区纯属浪费改小之后每个 Task 处理的数据量大一点、分区数少一点。如果你发现清洗作业在本地卡了很久先看这个参数。4.2 数据分组分析与可视化前的落盘格式选择清洗完的数据下一步就是分析。网约车场景里最经典的需求是“统计每个城市每个时段的订单量和订单总额”这个需求用 Spark SQL 和 DataFrame API 各有各的写法我习惯用 DataFrame API因为写起来像流水线result clean_df \ .withColumn(hour, hour(col(order_time))) \ .groupBy(city_id, hour) \ .agg( count(order_id).alias(order_cnt), sum(amount).alias(total_amount) ) \ .orderBy(col(city_id), col(hour)) result.write.mode(overwrite) \ .option(header, true) \ .csv(/data/result_csv)hour()函数从时间列里抽出小时数groupBy(city_id, hour)按城市和小时分组agg里count和sum分别算指标。最后落盘我写了 CSV因为要给后续的数据可视化用Flask 后端直接读 CSV 再返回给 ECharts 画柱状图和折线图是最短路径。如果这个结果还要继续参与计算比如和其他表 join就应该输出 parquet但如果只是给图表看CSV 反而更方便因为 BI 工具和前端解析 CSV 几乎零成本。这里要注意如果你在结果表上还留了 Spark 的 DataFrame 对象又反复去result.collect()driver 端内存可能会爆。可视化的数据量通常不大collect()之后转成 Pandas 再返回是常见做法但海量结果不要这么干可以直接写文件让后端去读。学习路线的 Day6 很多案例都是这样落盘的你看笔记里那个“网约车大数据综合项目——数据分析 spark 数据可视化 FlaskECharts”的链路中间这一步的格式选对了后面就不会卡壳。4.3 Spark SQL 与 Hive 的关系及几个常见面试题Day7 开始出现 Hive 相关的内容。很多文档把 Spark SQL 说成是“内存版 Hive”这个说法不准确但方向对。Spark SQL 的语法高度兼容 Hive而且可以通过spark.sql(...)直接写 Hive SQL只需要把hive-site.xml放到 classpath并配置spark.sql.catalogImplementationhive就能通过 Spark 读写 Hive 表。面试里围绕这个点的题非常固定我在带人复习时会让对方背熟两题。第一题repartition和coalesce的区别repartition会做全量 Shuffle把数据重新打散成指定分区数可以增也可以减coalesce默认不做 Shuffle只把相邻分区合并适合减少分区号但可能造成数据倾斜。第二题spark.sql.shuffle.partitions对什么生效它只影响 Spark SQL 执行 join、groupBy、distinct 时产生的 Shuffle 分区数不影响 RDD 的repartition后者要用spark.default.parallelism控制。这两个题背后的知识其实就是前面讲的宽依赖。面试官问分区相关的问题本质是想确认你懂不懂 Shuffle 的成本。你能答出“分区数不是越大越好”这一点比背一堆 API 名有用得多。5. Spark 3.0 入门必踩的坑现象、原因、解决5.1 本地能跑、提交集群就翻车序列化与依赖缺失现象代码在 spark-shell 里怎么跑都正常打包成 jar 提交到集群后报ClassNotFoundException或者NotSerializableException。原因集群的 Executor 进程拿不到你本地 IDE 里那份依赖。要么是自定义类没有实现Serializable闭包里的对象在 Shuffle 时要被序列化传输要么是第三方库没打进 jar提交时没有用--packages或者没有打 fat jar。解决自定义的 case class 或 bean 类先extends Serializable再说。依赖问题看报错里缺的是哪个类如果是 Spark 自身类检查pom.xml里 scope 是否错写成了 provided如果是第三方库用assembly插件打 fat jar或者直接用spark-submit --packages把依赖传进去。5.2 Executor 反复被 Kill内存参数永远要留出 overhead现象作业跑一半卡住日志里出现Container killed by YARN for exceeding memory limits然后 Executor 一个接一个挂掉。原因你只设置了--executor-memory 4g但 Spark 的堆外内存、线程栈、网络缓冲还要额外占用。这部分内存叫 overhead默认是executor-memory的 10%集群资源紧的时候很容易超。解决显式配置--conf spark.executor.memoryOverhead1g给堆外内存一个明确额度。如果作业本身有大量广播变量或 Kryo 缓存这个值还要再往上调。一句话记忆executor-memory管堆内memoryOverhead管堆外两个都别省。5.3 数据倾斜同一个 Stage 里某个 Task 永远跑不完现象Web UI 的 Stage 页面上大部分 Task 十几秒跑完某个 Task 要跑半小时以上而且它处理的输入数据量是其他 Task 的几十倍。原因groupBy 或 join 时 key 的分布不均匀比如某个城市 ID 占了大头或者 null 全被分到一个 Task。算力没变数据全堆在一处自然慢。解决第一选择是从业务上拆分把高频 key 单独拎出来处理第二选择是做加盐给 key 拼一个随机后缀拆成多个分区聚合完再合并。空值的处理要给结论如果空值业务上无意义直接filter过滤如果有意义用“null 随机数”打散。数据倾斜是大数据入门里最容易放弃治疗的问题但你要知道它能被解决——靠的是对宽依赖的理解而不是加大资源。5.4 结果文件碎成上千个小文件现象写出的结果目录里躺着几百上千个小文件每个只有几 KB。下游用 Hive 查的时候又慢又卡HDFS NameNode 压力也大。原因Spark 的 Shuffle 分区数决定了最终文件的个数默认 200但很多作业最终写出的数据量根本不够塞满 200 个分区。大量空分区和小分区直接落成空文件或碎片文件。解决写在之前加一个coalesce(n)或者repartition(n)n 数据量除以期望的单文件大小。举例1GB 数据想要 128MB 左右的文件coalesce(8)就够。Spark 3.0 也可以直接开spark.sql.adaptive.coalescePartitions.enabledtrue让 AQE 在运行时自动把分区数压到合理值。5.5 开了 AQE 后作业行为突变现象跟着教程把spark.sql.adaptive.enabled打开改了几个参数原本能跑的 join 突然变慢甚至结果出现奇怪的变化。原因AQE 是运行时优化它会在 Shuffle 结束后重新计算分区数和 join 策略。如果业务 SQL 本身写得不稳比如 join 条件里带函数、数据量估算偏差大AQE 可能做出和预期不同的决策。它确实是 3.0 的亮点但不是一个开关就能无脑开的东西。解决先小数据、单 Stage 验证把 AQE 的相关参数如spark.sql.adaptive.coalescePartitions.parallelismFirst调成 false 或 true 各跑一遍对比 Web UI 里的执行计划差异。如果 AQE 反而让作业变差关闭它回到静态执行计划也是一种正当选择。这个坑尤其误导入门者——教程说 3.0 默认开、性能更好但你没有意识到“性能更好”的前提是作业模型兼容 AQE。6. Day8 之后再往前一步用 Web UI 判断作业是否健康学完包里最后一天的东西你能提交作业、会看日志已经算入门了。但我想让你养成一个习惯每次作业跑完别急着关界面去:4040看一眼 Web UI 的四个关键位置。第一个是 Stages 列表。如果一个作业的 Stage 数量远多于算子数量说明有大量隐藏的宽依赖你该回去压缩 Shuffle。第二个是每个 Stage 的 Shuffle Read/Write 量。Write 巨大说明 map 端输出太多Read 巨大说明下游拉取不均衡这是数据倾斜的前兆。第三个是 Executors 页面里的Gc Time一列。如果某个 Executor 的 GC 时间占总运行时间 20% 以上优先把spark.memory.fraction从默认的 0.6 调低或者检查是不是collect了不该 collect 的数据。第四个是每个 Task 的 Duration 分布差异悬殊就是倾斜的直接证据。我自己在给新人验收时只问一个问题你的作业能不能不问别人仅靠 Web UI 说出瓶颈在哪。能说清楚的人对这个包的理解已经远超入门水平。这套看板排查习惯比记住任何调优参数都值钱因为它让你摆脱了把性能问题当玄学猜的状态。我现在的习惯是提交每个作业前先想清楚这个作业最怕的是什么。单纯跑数怕的是数据倾斜要写结果表怕的是小文件凌晨调度怕的是资源被抢。把这些预判写在笔记里每次跑完对照 Web UI 看是不是一一命中命中得越多你越了解自己的作业。希望帮到你。本文还有配套的精品资源点击获取
返回列表