Spark

Spark 是一个基于内存的分布式计算引擎。相比于Hadoop自带的MapReduce计算框架,Spark优势明显。 Spark一方面提供了更加灵活丰富的数据操作方式,有些需要分解成几轮 MapReduce作业的操作,可以在 Spark里一轮实现; 另一方面,每轮的计算结果都可以分布式地存放在内存中,下一轮作业直接从内存中读取数据,节省大量磁盘IO开销。

一、运行模式
1、基本概念
1.1 Yarn: Hadoop 资源管理系统,为应用提供统一的资源管理和调度。
  • ResourceManger:处理客户端请求、启动和监控ApplicationMaster、监控NodeManager、资源的分配与调度
  • NodeManger:管理单个节点上的资源、处理来自ResourceManager的命令、处理来自ApplicationMaster的命令
  • ApplicationMaster:为应用程序申请资源
1.2 Spark: 分布式计算框架
  • Application:用户自己写的Spark应用程序,批处理作业的集合。Application的main方法为应用程序的入口,用户通过Spark的API,定义了RDD和对RDD的操作
  • SparkContext:Spark最重要的API,用户逻辑与Spark集群主要的交互接口。
  • Driver:主控进程。SparkContext运行在其中,负责产生DAG,提交Job,转化Task。
  • Executor:负责执行Task,并将结果返回给Driver,同时也提供缓存的RDD的功能
2、两种运行模式
  • yarn-client模式driver端在本地,可以与集群进行调度和通讯,进行交互式作业。但如果作业很多则会造成client端压力很大。
  • yarn-cluster模式driver运行在AM中,client提交完作业就可以关掉了,但无法进行交互式作业。
二、一个SQL语句是如何被转换的?
1. 解析阶段 (Parsing)
  • 词法与语法分析:使用 ANTLR 等工具将输入的 SQL 字符串转换成抽象语法树(AST,即未解析的逻辑计划 Unresolved Logical Plan)。
  • 检查元数据:利用 SessionCatalog 结合元数据信息(如 Hive MetaStore 或内存中的表结构),对 AST 中的表名、字段名、函数名进行绑定和有效性校验,生成已解析的逻辑计划(Analyzed Logical Plan)。
2. 优化阶段 (Logical Optimization)
  • 规则优化:Catalyst 优化器内置了大量的优化规则(Rule-based Optimization),如谓词下推(Predicate Pushdown)、列裁剪(Column Pruning)、常量折叠(Constant Folding)等。

image.png

  • 代价优化:基于统计信息(Cost-based Optimization, CBO)评估不同的执行路径,调整连接顺序或选择更优的 Join 策略,最终输出优化的逻辑计划(Optimized Logical Plan)。
3. 物理计划阶段 (Physical Planning)
  • 生成物理计划:将优化后的逻辑计划转换成一个或多个物理计划(Physical Plan),并通过贪心算法或 Cost 模型选择代价最低的物理执行策略(例如选择 BroadcastHashJoin 还是 SortMergeJoin)。
  • 准备执行:对物理算子进行进一步处理,比如插入必要的 Shuffle 算子(Exchange)或行列转换操作。

Exchange算子是实现数据并行化的重要算子,用于解决数据分布(Distribution)相关问题。继承它的子类有两种:

- BroadcastExchange:广播操作。

- ShuffleExchange:通过shuffle进行重分区。

4. 代码生成阶段 (Code Generation)
  • 生成 Java 字节码:利用 Spark 的 Whole-Stage Code Generation(全阶段代码生成)技术,将整个物理计划编译成高效的 Java 字节码。
  • 转为 RDD 执行:把编译好的代码包装在 RDD 的算子中,直接提交给底层的 Spark Core 调度运行,从而避开传统的解释执行开销,大幅提升计算性能。
三、核心抽象RDD(Resilient Distributed DateSet)

RDD(弹性分布式数据集) 是 Spark 的基石,可以理解为一个分布在多台机器上的、只读的数据集合

1、RDD特性:
  • 分区(Partition):数据被切成多个分区散落在不同节点,分区是并行计算的最小单位。
  • 不可变:RDD 一旦创建就不能修改,任何转换都是生成一个新的 RDD。
  • 血缘(Lineage):每个 RDD 都记得自己是怎么从上一个 RDD 算出来的。某个分区的数据丢了,不需要重算全部,顺着血缘把那个分区重新算一遍就行——这就是「弹性」的含义,也是 Spark 容错的核心机制。
2、RDD 操作

1. Transformation 算子:通常用来将RDD内的元素进行转换,并返回一个新的RDD,它们只搭建管道,不执行计算

  • map(func):对 RDD 中的每个元素应用一次函数,返回新 RDD。
  • filter(func):过滤掉不符合条件的元素,返回新 RDD。
  • flatMap(func):类似 map,但每个输入元素可以映射为 0 或多个输出元素。
  • reduceByKey(func):针对 Key-Value 类型的数据,按照相同的 Key 来进行聚合。
  • groupByKey():将 Key-Value 类型的数据按照 Key 进行分组。

2. Action 算子(行动):它们会触发真正的任务调度与运算:

  • collect():将 RDD 的所有元素拉取到 Driver 端内存中,形成一个普通集合。
  • count():返回 RDD 中元素的总个数。
  • reduce(func):在驱动程序中对数据集中的元素进行聚合计算。
  • first():返回 RDD 中的第一个元素。
  • saveAsTextFile(path):将数据集中的元素保存到文本文件或 HDFS 中。
四、常见的概念及其关系
1、Application、Job、Stage、Task之间的关系
text
1Application Spark
2 Job Action Job
3 Stage
4 Task Task

Stage 怎么切分,取决于依赖类型

  • 窄依赖:父 RDD 的每个分区最多被一个子分区使用(如 mapfilter)。不需要跨节点传数据,可以在同一个 Stage 里连续执行。
  • 宽依赖:父 RDD 的一个分区被多个子分区使用(如 groupByKeyjoin)。必须做 Shuffle,因此成为 Stage 的分界线。

所以一句话:遇到宽依赖就断开,切出一个新 Stage。Stage 越多意味着 Shuffle 越多,调优的方向往往就是想办法减少宽依赖。

2、Driver、Executor、Core、Task、队列之间的关系

Driver 是 Spark 应用的大脑。它负责:

  • 解析 SQL / 代码,生成执行计划
  • 划分 Job、Stage、Task
  • 向集群申请资源
  • 调度 Task 到 Executor 上执行
  • 收集任务状态和部分结果

Executor 是真正干活的进程。它负责:

  • 执行 Driver 分配过来的 Task
  • 缓存数据
  • Shuffle 读写
  • 向 Driver 汇报状态

Core 决定 Executor 的并发能力。一般来说:1 个 core 同一时间通常跑 1 个 task,例如:

  • executor-cores = 4
  • num-executors = 10
  • 那么总并发 Task 数大约是:4 \* 10 = 40

Task 是 Spark 最小执行单元。

  • 一个 Stage 会被拆成多个 Task,通常和分区数对应。
  • 一个 partition 通常对应一个 task,如果某个 Stage 有 200 个分区,就会有约 200 个 Task。

队列 是集群资源池,比如 Yarn 队列。队列决定你的 Spark 应用最多能拿到多少:

  • CPU
  • Memory
  • Executor 数
  • 并发资源

所以队列是上限,Spark 参数是在这个上限内分配资源。

3、CPU、Memory和Spark参数的对应关系

CPU ≈ executor 数 \* executor cores

Memory ≈ executor 数 \* executor memory

并发 task 数 ≈ executor 数 \* executor cores

4、常用配置参数
Driver:主控进程。SparkContext运行在其中,负责产生DAG,提交Job调度task。task数量越多,消耗内存和网络io越多。
  • spark.driver.memory=12G:driver的堆内存
  • spark.driver.cores:driver的core数
  • spark.driver.memoryOverhead=6144:driver的堆外内存
Executor::负责执行Task,并将结果返回给Driver,同时也提供缓存的RDD的功能。shuffle越多或数据量越大,消耗的内存越多。task计算的并行度由core个数决定。core越多,计算越快,task所占用的内存空间就越小。所以这是一个trade-off,如果任务的shuffle量很小,主要是一些计算、过滤操作,那么可以多核少内存,甚至可以开启spark.vcore.boost.ratio=2(原spark.vcore.boost=2),一个core来跑两个task,提高cpu利用率。如果shuffle很重,那么就需要看shuffle write/read 指标来调节内存大小。
  • spark.executor.memory=8G:每个executor的堆内存大小。风神默认5g。
  • spark.executor.cores=4:每个executor的核数,默认1个core运行1个task。风神默认3个
  • spark.executor.memoryOverhead=6144:每个executor的堆外内存,默认取max(384MB,spark.executor.memory \*0.1),dorado默认6G,风神默认3G
Executor个数:建议使用动态资源分配(默认开启)。使它可以根据工作负载动态调整应用程序占用的资源。这意味着,如果不再使用资源,应用程序可能会将资源返回给集群,并在稍后需要时再次请求资源。如果多个应用程序共享Spark集群中的资源,该特性尤其有用。Driver会判断如果没有在task pending的情况,资源会被释放掉
  • spark.dynamicAllocation.enabled=true,开启参数
  • spark.dynamicAllocation.minExecutors=5,executor最小申请个数
  • spark.dynamicAllocation.maxExecutors=900,executor最大申请个数
  • spark.dynamicAllocation.initialExecutors=5,初始申请executor个数,默认等于minExecutors
  • spark.dynamicAllocation.executorIdleTimeout=120s,一个executor空闲超过该参数时,自动释放资源
  • spark.dynamicAllocation.schedulerBacklogTimeout(默认1秒),如果有task pending超过该参数,开启资源申请
  • spark.dynamicAllocation.sustainedSchedulerBacklogTimeout(默认等于上面的参数),当pending task存在以后,每隔该参数进行一个资源的申请(如果资源有,会以1、2、4、8指数的申请)
  • spark.shuffle.service.enabled=true,在不删除由executors产生的shuffle文件的情况下删除executors
五、Spark Web UI
1、Jobs

在浏览器中打开tracking URL后,默认进入Jobs页。Jobs展示的是整个spark应用任务的job整体信息:

  • User: spark任务提交的用户,用以进行权限控制与资源分配。
  • Total Uptime: spark application总的运行时间,从appmaster开始运行到结束的整体时间。
  • Scheduling Mode: application中task任务的调度策略,由参数spark.scheduler.mode来设置,可选的参数有FAIR和FIFO,我司默认是FAIR。这与yarn的资源调度策略的层级不同,yarn的资源调度是针对集群中不同application间的,而spark scheduler mode则是针对application内部task set级别的资源分配。yarn 的FAIR模型是一种公平模型,相当于每个任务轮换使用资源等,这样能使的小job能很快执行,而不用等大job完成才执行了。
  • Completed Jobs: 已完成Job的基本信息,如想查看某一个Job的详细情况,可点击对应Job进行查看。
  • Active Jobs: 正在运行的Job的基本信息。
  • Event Timeline: 在application应用运行期间,Job和Exector的增加和删除事件进行图形化的展现。这个就是用来表示调度job何时启动何时结束,以及Excutor何时加入何时移除。我们可以很方便看到哪些job已经运行完成,使用了多少Excutor,哪些正在运行。

Job默认都是串行提交运行的,如果Job间没有依赖,可以使用多线程并行提交Job,实现Job并发。

image.png

2、Jobs Detail

在Jobs页面点击进入某个Job之后,可以查看某一Job的详细信息:

  • Status: 展示Job的当前状态信息。
  • Active Stages: 正在运行的stages信息,点击某个stage可进入查看具体的stage信息。
  • Pending Stages: 排队的stages信息,根据解析的DAG图stage可并发提交运行,而有依赖的stage未运行完时则处于等待队列中。
  • Completed Stages: 已经完成的stages信息。
  • Event Timeline: 展示当前Job运行期间stage的提交与结束、Executor的加入与退出等事件信息。
  • DAG Visualization: 当前Job所包含的所有stage信息(stage中包含的明细的tranformation操作),以及各stage间的DAG依赖图。DAG也是一种调度模型,在spark的作业调度中,有很多作业存在依赖关系,所以没有依赖关系的作业可以并行执行,有依赖的作业不能并行执行

image.png

3、Stages

在Job Detail页点击进入某个stage后,可以查看某一stage的详细信息:

  • Total time across all tasks: 当前stage中所有task花费的时间和。
  • Locality Level Summary: 不同本地化级别下的任务数,本地化级别是指数据与计算间的关系。
  • Input Size/Records: 输入的数据字节数大小/记录条数。这里可以通过参数控制任务初始并行度。
  • Output Size / Records: 输出的数据字节数大小/记录条数。
  • Shuffle Write: 为下一个依赖的stage提供输入数据,shuffle过程中通过网络传输的数据字节数/记录条数。应该尽量减少shuffle的数据量及其操作次数,这是spark任务优化的一条基本原则。
  • Shuffle read:总的shuffle字节数,包括本地节点和远程节点的数据
  • DAG Visualization: 当前stage中包含的详细的tranformation操作流程图。
  • Show Additional Metrics:显示额外的一些指标。鼠标移上去会有相应的解释。
  • Event Timeline: 清楚地展示在每个Executor上各个task的各个阶段的时间统计信息,可以清楚地看到task任务时间是否有明显倾斜,以及倾斜的时间主要是属于哪个阶段,从而有针对性的进行优化。
  • Summary Metrics for XXX Completed Tasks:已完成的Task的指标摘要。
  • Duration:task持续时间
  • GC Time:task gc消耗时间
  • Output Write Time:输出写时间
  • Output Size/ Records : 输出数量大小,条数。
  • Shuffle HDFS Read Time: shuffle中间结果从hdfs读取时间
  • Shuffle Read Size/Records: shuffle 读入数据大小/条数
  • Shuffle Spill Time: shuffle 中间结果溢写时间。
  • Shuffle spill (memory):shuffle 溢写使用的内存大小。
  • Shuffle spill (disk):shuffle 溢写使用的硬盘大小。
  • Aggregated Metrics by Executor: 汇总指标。将task运行的指标信息按executor做聚合后的统计信息,并可查看某个Excutor上任务运行的日志信息。这里可以看到executor完成task的情况,读入的数据量,shuffle write、shuffle read 数据量,根据这些指标判断节点是否健康。
  • Tasks: 当前stage中所有任务运行的明细信息。右边部分都是需要关注的点。

image.pngimage.png

4、Storage

storage页面能看出application当前使用的缓存情况,可以看到有哪些RDD被缓存了,以及占用的内存资源。如果job在执行时持久化(persist)/缓存(cache)了一个RDD,那么RDD的信息可以在这个选项卡中查看。

image.png

5、Storage Detail

点击某个RDD即可查看该RDD缓存的详细信息,包括缓存在哪个Executor中,使用的block情况,RDD上分区(partitions)的信息以及存储RDD的主机的地址。

image.png

6、Environment

Environment选项卡提供有关Spark应用程序(或SparkContext)中使用的各种属性和环境变量的信息。用户可以通过这个选项卡得到非常有用的各种Spark属性信息,而不用去翻找属性配置文件。通过平台工具提交任务时,其默认参数配置可以在这里查找。

image.png

7. Executor

Executors选项卡提供了关于内存、CPU核和其他被Executors使用的资源的信息。这些信息在Executor级别和汇总级别都可以获取到。一方面通过它可以看出来每个excutor是否发生了数据倾斜,另一方面可以具体分析目前的应用是否产生了大量的shuffle,是否可以通过增加并行度来减少shuffle的数据量。

  • Summary: 该application运行过程中使用Executor的统计信息。
  • Executors: 每个Excutor的详细信息(包含driver),可以点击查看某个Executor中任务运行的详细日志。

image.png

六、Spark内存
  • spark.memory.fraction=0.6 统一内存占比。在shuffle很重、内存很大的情况下可以适当调大该参数。

image.png

评论 (0)

登录后参与评论。

还没有评论,来做第一个。

登录后可以选中正文添加批注(仅自己可见)。

Spark