ARTICLE DETAIL

资讯详情

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

Spark 3.0入门到精通:环境搭建、RDD、SparkSQL与调优实战

Spark 3.0入门到精通:环境搭建、RDD、SparkSQL与调优实战 简介面向大数据入门者与Spark自学者这份配套资料基于Spark 3.0.1稳定版按1-8天的学习节奏组织代码与笔记内容覆盖环境搭建、SparkCore、SparkStreaming、SparkSQL、StructuredStreaming、Spark综合案例、多语言开发、3.0新特性及性能调优等核心模块。课程从基础环境搭建逐步延伸到实时流处理与性能优化递进式安排便于读者按天推进。压缩包共244个文件大小86.9MB以PNG截图、Markdown笔记、Scala源码和ZIP压缩包为主另有JSON配置文件截图直观展示运行效果Markdown笔记承载每日知识点Scala源码供动手实践整体目录结构清晰能快速定位每日章节。已有609人下载学习适合希望系统迈入Spark门槛、边写边练的零基础或初级读者。每日笔记按天拆分配合大量运行截图与可复用代码既能帮助从零搭建Spark开发环境也能作为日常编码时的参考综合案例与调优章节进一步展示真实场景下的应用与性能优化思路有助于提升动手实践能力。1. 为什么Spark 3.0值得学从“过时框架”的误解说起大数据开发这个岗位面试十次有八次绕不开Spark。但网上很多教程还停留在Spark 2.x甚至1.6的写法导致不少人学完发现连RDD的API都变了。这套2021贺岁版的Spark 3.0入门到精通资源用的是官方在2020年9月8日发布的Spark 3.0.1稳定版共9个章节覆盖环境搭建、SparkCore、SparkStreaming、SparkSQL、StructuredStreaming、综合案例、多语言开发、3.0新特性和性能调优配套1-8天的完整代码和Markdown笔记正好补齐了从零到上手开发这中间最缺的那段路。适合刚接触大数据、想快速把环境跑起来并且能看懂调优参数的新手也适合有一两年Hadoop经验、想从MapReduce思维切换到Spark算子的从业者。接下来我直接用这套资源里真实的目录和代码拆一遍完整的学习路径和落地时容易踩的坑。2. Spark 3.0.1环境搭建版本选型与三个核心配置点先说为什么必须锁版本。很多新手从官网直接下最新版结果装完发现配套的Hadoop版本对不上或者JDK版本太新导致启动报错。这套资源选Spark 3.0.1是有讲究的它和Hadoop 3.2、JDK 8是经过大量生产验证的组合网上遇到的问题也最多搜解决方案最容易。Spark 3.0相比2.x最大的变化是引入了动态分区裁剪、自适应查询执行AQE这些特性而这些特性在2.4及以前版本里要么没有、要么是实验性的所以学3.0.1不会学到过时的写法。2.1 前置环境JDK 8与Hadoop的关系Spark不强制依赖Hadoop的HDFS但Standalone模式下读取本地文件、测试WordCount时环境变量里必须配好JAVA_HOME。这套资源里没有内置JDK安装包Day01笔记里默认你已经配好了。我一般建议用/usr/lib/jvm/java-1.8.0-openjdk这个路径因为CentOS上直接用yum install java-1.8.0-openjdk安装就是它不会出现路径找不到的问题。安装完成后先验证java -version echo $JAVA_HOME如果JAVA_HOME是空的需要在/etc/profile里追加export JAVA_HOME/usr/lib/jvm/java-1.8.0-openjdk export PATH$PATH:$JAVA_HOME/bin然后source /etc/profile。注意这里有个分布式的坑如果你有三台机器每台都要执行同样的配置只改一台的JAVA_HOME会导致Worker节点启动失败。这个在Day01笔记里没有刻意强调但我自己第一次搭集群时就因为只配了master节点花了半小时排查为什么worker总是连不上。2.2 下载与解压两个容易出问题的点Spark 3.0.1的下载包有两种spark-3.0.1-bin-hadoop2.7.tgz和spark-3.0.1-bin-hadoop3.2.tgz。这套资源默认搭配Hadoop 3.2你如果用hadoop2.7的包去连HDFS会报Incompatible clusterIDs之类的协议错误。解压路径建议统一放在/opt/spark并且给目录一个不带版本号的软链接cd /opt tar -zxvf spark-3.0.1-bin-hadoop3.2.tgz ln -s spark-3.0.1-bin-hadoop3.2 spark这样做的原因是后面Spark依赖的SPARK_HOME指向/opt/spark下次升级版本只需要重新做一次软链接代码和配置不用动。视频里有一步会去改spark-env.sh里的JAVA_HOME我建议同时也把SPARK_MASTER_HOST设置成内网IP不要用localhost否则你用spark://localhost:7077提交任务时远程客户端连不上。2.3 spark-env.sh与spark-defaults.conf的核心参数这套资源在Day01里给出的spark-env.sh配置比较精简但实际生产环境有几个参数一定要在这个阶段就理解export SPARK_MASTER_HOST192.168.1.100 export SPARK_WORKER_CORES4 export SPARK_WORKER_MEMORY8g export SPARK_DRIVER_MEMORY2g export SPARK_EXECUTOR_MEMORY4gSPARK_WORKER_CORES是本机Worker能够分配给Executor的总核数SPARK_WORKER_MEMORY是Worker总内存。注意这两个值是上限不是每个任务独享的。新手最容易犯的错是把它们当成单个Executor的资源于是提交任务时--executor-memory 6g结果Worker总共才8g任务一直卡在Waiting for resources状态。spark-defaults.conf里要提前打开一个Spark 3.0的关键开关spark.sql.adaptive.enabled true spark.sql.adaptive.coalescePartitions.enabled true这就是自适应查询执行的开关3.0.1里AQE还属于稳定功能打开后能自动合并shuffle产生的小文件分区减少task数量对后续SparkSQL的查询性能有明显的提升。如果你用默认关闭状态跑完整个课程后面调优章节有些实验效果会看不出来。启动并验证环境/opt/spark/sbin/start-master.sh /opt/spark/sbin/start-worker.sh spark://192.168.1.100:7077 /opt/spark/bin/spark-shell --master spark://192.168.1.100:7077看到Spark context Web UI available at http://192.168.1.100:4040就说明环境正常了。我遇到一个很玄学的问题是防火墙开着导致4040端口能访问、7077端口连不上排查方法是telnet 192.168.1.100 7077如果超时就systemctl stop firewalld再试确认是防火墙问题再放行端口不迟。3. RDD编程与SparkCoreDay01到Day03笔记里的核心抽象环境搭完接下来进入这套资源里最重的部分SparkCore。Day01到Day03笔记基本围绕RDD展开从创建、算子使用到持久化策略。RDD这个抽象理解不到位后面的SparkSQL和StructuredStreaming学起来全是空中楼阁因为很多性能问题要回到RDD的调度层面去找原因。3.1 RDD的创建parallelize与textFile的适用边界Day01笔记里给的创建方式是parallelize和textFile两种。但代码很容易被照抄后出问题val rdd1 sc.parallelize(List(1, 2, 3, 4, 5), 2) val rdd2 sc.textFile(hdfs://master:9000/input/wordcount.txt, 3)parallelize的第二个参数是分区数对应的是内存中集合的分片数量不是数据文件的数量。而textFile的第二个参数是最小分区数实际分区数由HDFS的block数量决定minPartitions只是个下限。新手容易犯的错是觉得分区数自己定了就一定会按这个数执行其实当文件大小超过minPartitions * blockSize时Spark会按block数切更多分区。还有一点在Day01里容易忽略textFile读取目录时只读文件不读子目录如果数据是按天分目录存的要读多天数据得把多个路径用逗号拼起来val rdd sc.textFile(/data/day01,/data/day02,/data/day03)这条路如果走通了后面综合案例里处理多天日志就会很顺畅。3.2 Transformation与Action惰性求值如何影响你的代码Day02笔记里强调最多的是惰性求值。map、filter、flatMap这些算子只是构建DAG不会真正计算遇到collect、count、saveAsTextFile这些Action才会触发job。这个机制带来的第一个麻烦是你在map操作里写的println不会执行。val rdd sc.textFile(/opt/spark/README.md) val words rdd.flatMap(_.split( )).filter(_.length 3) words.map(word println(sprocessing: $word))这段代码运行后控制台什么都打印不出来因为map只是记录了函数没有真正运行。要看到中间结果必须触发Actionwords.collect().foreach(println)但collect是把所有数据拉到Driver端数据量大时直接OOM。根本原因是惰性求值的设计意图是减少IO和计算量不是方便你调试。我一般调试时会用take(10)替代collect既能触发计算又不会拉全量数据。这个习惯在Day03的shuffle调优章节里还会用到。3.3 持久化与Checkpoint两个必须养成的习惯Day03笔记里cache、persist、checkpoint三者的区别讲得很细。但它的实际价值要放在迭代计算或shuffle重算场景里才有体感。下面这个场景在综合案例里比较典型val logs sc.textFile(/data/access.log).map(parseLog) logs.cache() val totalCount logs.count() val errorCount logs.filter(_.status 500).count() val urlCount logs.map(_.url).distinct().count()三个Action都从logs这个RDD开始如果没有cache()每个Action都会从源头textFile开始重新计算解析逻辑三次Action等于三倍读取和解析开销。cache()在Spark 3.0里的默认行为是MEMORY_ONLY如果数据放不下多余的分区会重新计算而不是溢写到磁盘这个在Day03笔记里没有明说但实际调优时很重要。Checkpoint是另一种持久化思维它把RDD数据写到类似HDFS的可靠存储并且会切断RDD的血统链。两者的本质区别是cache保留血统而checkpoint不保留所以checkpoint的恢复不会因为某个中间算子失败而重算整个链条。但checkpoint的代价是高能不用就不用千万别每步都打。4. SparkSQL与Streaming从结构化数据到实时计算的衔接SparkSQL和Streaming是这套资源里Workload比较重的部分分别对应Day04、Day05、Day06。SparkSQL的重点在DataFrame/Dataset的统一APIStreaming部分的重点则是SparkStreaming和StructuredStreaming两种模型的差异。这两个模块放在一起学是因为它们共享Spark SQL引擎的Catalyst优化器很多性能调参方式是一样的。4.1 DataFrame与Dataset代码里怎么选、为什么Day04笔记里大量使用spark.read.option(header, true).csv(...)来动态推断Schema但真正要落地时我建议显式指定Schema而不是依赖推断因为推断的Schema经常把数值列识别成StringTypefrom pyspark.sql.types import StructType, StructField, StringType, IntegerType, LongType schema StructType([ StructField(order_id, StringType(), True), StructField(user_id, StringType(), True), StructField(amount, IntegerType(), True), StructField(ts, LongType(), True) ]) df spark.read \ .option(header, true) \ .schema(schema) \ .csv(/data/orders.csv)显式定义Schema有三个直接好处避免CSV首行被当成表头、避免数值列被推断成String导致后续聚合结果错乱、避免每次读取重复做类型推断耗时。注意StructField里的第三个参数是nullable这个字段在Spark 3.0里如果不写会默认为true但在生产环境如果明确某列有值写成false可以帮助优化器在谓词下推时更激进。Dataset在Python API里没有完全对应的概念PySpark的DataFrame就是Dataset[Row]。这套资源主要以Scala代码为主如果你用PySpark跑ds.map()这种强类型算子是走不了的要把思路切换成select、withColumn这些表达式风格。4.2 SparkStreaming与StructuredStreaming的取舍Day05和Day06分别讲了SparkStreaming和StructuredStreaming。一个容易被忽略的关键点是SparkStreaming基于RDD的微批模型在Spark 3.0里已经进入维护模式而StructuredStreaming在3.0中作为新引擎才是推荐的流处理方案。如果你从Day05就开始学要记住它的思维模型和StructuredStreaming不完全一致。StructuredStreaming最核心的是基于事件时间的窗口聚合val input spark.readStream .format(kafka) .option(kafka.bootstrap.servers, master:9092,worker01:9092) .option(subscribe, orders) .load() val query input .selectExpr(CAST(value AS STRING)) .as[String] .map(parseOrder) .withWatermark(eventTime, 5 minutes) .groupBy(window($eventTime, 10 minutes, 5 minutes)) .agg(sum($amount).as(total_amount)) .writeStream .outputMode(append) .format(console) .start().withWatermark(eventTime, 5 minutes)是处理延迟数据的关键5分钟是允许数据迟到多少时间后仍参与窗口计算。.groupBy(window($eventTime, 10 minutes, 5 minutes))的第一个参数是窗口大小第二个是滑动步长。注意事件时间必须是从数据里解析出来的字段不能是current_timestamp()否则所有数据都落在当前窗口里等于没做窗口聚合。实际测试时Kafka如果还没装可以把format改成socket并配合nc -lk 9999手动发数据Day06里应该有类似的例子。但用socket模式时withWatermark是不生效的因为没有内置的eventTime字段你需要自己在输入字符串里拼接时间戳。5. Spark避坑指南从环境变量到内存溢出的五个高频问题这套资源从Day01到Day06我按自己跑通的经验整理了五个最容易卡住的坑每一条都是实际测试中真实遇到过的不是理论推演。坑一Worker启动成功但Executor总是启动失败。现象是Master Web UI里能看到Worker节点但提交任务后任务卡在WAITING状态日志显示Executor启动失败。原因是SPARK_WORKER_MEMORY和spark.executor.memory设置不匹配Executor申请的内存大于Worker剩余可用内存。排查方法是看/opt/spark/logs/spark-org.apache.spark.deploy.worker-*.out如果提示Failed to allocate memory把spark-submit里的--executor-memory降到Worker内存的一半以下再试。我在实际部署时习惯把SPARK_WORKER_MEMORY设成物理机的75%留出系统余量Executormemory不超过2g这样很少触发容器分配失败。坑二sc.textFile读取本地文件报FileNotFound。现象是在Spark Shell里执行sc.textFile(/home/user/data.txt)Driver端报文件找不到但ls明明能看到文件。原因是Spark默认使用HDFS协议解析路径本地文件必须显式加file://前缀sc.textFile(file:///home/user/data.txt)。如果是Standalone集群模式Worker节点的路径必须和Driver端一致把文件只放到master上是没用的每个Worker都要有这份文件否则执行collect时报NoSuchFileException的却是Worker端。坑三SparkSQL读取CSV时列数对不上。现象是df.show()能正常显示但df.count()报错或者某些行数据全变成null。根本原因是CSV里某些行的字段被逗号分隔后数量不一致常见于字段内容里包含逗号但没加引号。解决办法是读取时加个modedf spark.read.option(header, true).option(mode, PERMISSIVE).csv(/data/messy.csv)PERMISSIVE是默认模式它会吞掉解析错误并把这行对应的字段置为null还有DROPMALFORMED直接丢弃坏行、FAILFAST遇到坏行直接报错。排查时先确定自己的数据属于哪种脏数据再选对应的模式。如果不需要这些脏行参与统计DROPMALFORMED比事后过滤更高效。坑四Kafka与SparkStreaming的版本不匹配。现象是消费Kafka数据时一直报OffsetOutOfRangeException或者启动时提示kafka.clients.consumer.internals相关的NoSuchMethodError。原因是spark-streaming-kafka-0-10这个依赖的Kafka client版本和实际部署的Kafka broker不一致。我一般显式在pom里指定kafka-clients与spark-streaming-kafka-0-10的版本配套比如kafka_2.12-2.5.0对应spark-streaming-kafka-0-10_2.12:3.0.1不要再引入额外的kafka-clients依赖让Spark传递依赖自己解析版本。坑五groupBy后直接collect导致Driver OOM。现象是数据体量不大但collect后内存溢出。原因是groupBy产生宽依赖shuffle后每个key对应的value集合在Driver端被组装成一次性的大对象这个对象的大小远超原始单条数据。Spark 3.0里可以用toDF的orderBy或者agg配合collect_list避免直接拉全量更稳妥的做法是把聚合结果直接落地到存储再读df.groupBy(category) .agg(sum(amount).as(total)) .write.mode(overwrite).parquet(/output/category_amount)落地到Parquet后可以增量查看不会因为一次性加载所有聚合数据而压垮Driver。6. 综合案例与多语言开发把离线任务跑成闭环的关键技巧整套资源最后指向两个文件Spark综合案例.md和Spark-多语言开发.md。综合案例在这套资源里已经把几个核心组件串起来了从数据清洗到聚合统计最后落结果到存储。多语言开发则是用Scala之外的语言做Spark开发实际面试和工作中用到的是Java和Python两种方向。6.1 清洗脏数据用RDD还是DataFrame综合案例第一部分基本是数据清洗。我复现时发现用DataFrame的filter和withColumn做这套清洗比纯RDD的map方式代码量少一半以上。比如去掉字段数为空的记录DataFrame写法是df_cleaned df.filter(df[user_id].isNotNull() (df[amount] 0))而RDD要逐行._1取字段再判断可读性和维护成本都差一截。数据清洗后用df_cleaned.createOrReplaceTempView(orders)注册临时表再用SQL写统计逻辑语序比冗长的DataFrame链式调用直观很多。6.2 多语言开发PySpark与Scala代码的翻译对照Day07里应该给了同一套逻辑的Scala和Python版本对比。核心要掌握的是PySpark的DataFrame语法和Scala几乎一一对应但有些算子有差异。比如Scala里rdd.map(x x * 2)在PySpark里要写rdd.map(lambda x: x * 2)Scala的$amount列引用在PySpark里换成col(amount)。翻译一套简单逻辑建议直接从官方文档的示例入手别自己瞎猜。6.3 快速定位问题日志检索的三个级别任务失败时先看三个地方stderr日志里的Exception、Web UI里对应Job的DAG Visualization、以及Executor日志中的Task失败详情。我习惯在spark-submit后面加--conf spark.log.levelINFO但排查时临时改WARN也能减少刷屏。从那以后我每次提交生产任务都强制走一遍先看日志关键字、再看Web UI的Stage耗时、最后查Executor状态这套流程省下来的时间比什么参数调优都多。希望帮到你。本文还有配套的精品资源点击获取
返回列表