spark20问-1
1、Spark 是什么?它与 Hadoop 有什么区别?
spark 是一个分布式大数据处理框架,提供了内存计算的能力。
相比于 Hadoop:
spark 以内存计算为核心,避免了大量的 I/O,Hadoop 则是采用磁盘计算模式,数据存储、计算过程需要多次读写磁盘,速度较慢;
spark 具备独立 DAG 任务调度引擎,可以优化任务执行顺序,灵活调度资源,Hadoop 由 HDFS 和 MapReduce 构成;
spark 提供 RDD 编程模型、Dataset、DataFrame 高层次抽象接口,支持 java、python、scala、r,Hadoop 以 MapReduce 键值对编程模型为主,多数情况下需要手动处理复杂任务的拆分和合并;
spark 适用于机器学习、离线数据处理、交互式查询;Hadoop 适用于对数据持久化要求高的大规模离线的批处理任务
2、在 Spark 中,什么是 RDD?它的特点是什么?
RDD 是分布式数据集合,提供一种容错、高效、并行的数据处理方式;
RDD 有以下几个特点:
a. 弹性 Resilient:记录生成每一个 RDD 的一系列操作的元数据 Lineage,节点失效时,RDD 可以根据这个元数据重新计算丢失的数据
b. 分布式 Distributed:RDD 是一个分布在集群节点上的数据集合,可跨越多个节点进行存储和操作,从而充分利用集群的计算和存储资源
c. 不可变:RDD 不可变,对 RDD 的每一个操作都会生成一个新的 RDD,不会改变原始的 RDD
d. 惰性计算:RDD 所有操作都是惰性的,不会立即执行,而是先记录下来,当遇到动作操作如 collect / save 等时,才会触发计算
# 创建 rdd
val rddFromFile = sc.textFile("hdfs://path/to/file")
val rddFromParallelize = sc.parallelize(Seq(1, 2, 3, 4, 5))
# rdd 转换操作
val rdd1 = sc.parallelize(Seq(1, 2, 3))
val rdd2 = rdd1.map(_ * 2)
val rdd3 = rdd2.filter(_ > 2)
# rdd 动作操作
val data = rdd3.collect()
val count = rdd3.count()
rdd3.saveAsTextFile("hdfs://path/to/output")
# 持久化
rdd3.persist(StorageLevel.MEMORY_ONLY)
3、Spark 的基本架构是什么?主要包括哪些组件?
Driver:执行用户 main 函数,将代码转化为 Task 交给集群执行。Driver 维护了 包括 RDD图 在内的 Spark 应用程序的执行状态和所有结构信息
Cluster Manager:负责整个集群的资源分配,决定哪些计算资源可以用于运行任务,有 Hadoop YARN、Apache Mesos、Spark Standalone Manager。
Worker Nodes:实际执行任务的机器,每一个节点上都有一个 Worker 进程,它们为 spark 提供计算和内存资源,用于运行 Executor。
Executor:spark 应用程序在 Worker 节点上执行时创建的计算资源实例,每个应用程序在启动时会获得一个/多个 Executor 用于执行任务进行数据存储。
RDD:可分区、容错、只读、并行执行的 数据集合。后续被 DataFrame 和 Dataset 逐步替代。
4、在 Spark 中,如何创建一个 RDD?
有两种:1、从文件或外部数据源创建 2、从已有集合创建
val sc = new SparkContext(conf)
val rddFromFile = sc.textFile("hdfs://path/to/file.txt")
val sc = new SparkContext(conf)
val rddFromCollection = sc.parallelize(Seq(1, 2, 3, 4, 5))
# 3 种 rdd 的转化操作:
# map: 对 rdd 中每一个元素进行一次操作,并返回一个新的 rdd
val rddMapped = rddFromCollection.map(x => x * 2)
# filter:过滤掉不满足条件的元素
val rddFiltered = rddFromCollection.filter(x => x > 3)
# flatMap:与 map 相似,但可以返回一个包含多个元素的集合
val rddFlatMapped = rddFromCollection.flatMap(x => Seq(x, x * 2))
# 3 种 rdd 的行动操作
# collect:收集 rdd 中的所有元素,返回给驱动程序
val collectedData = rddFromCollection.collect()
# count:返回 rdd 中元素的数量
val count = rddFromCollection.count()
# reduce:使用用户提供的二元操作函数来合并 rdd 中的元素
val sum = rddFromCollection.reduce(_ + _)
# 持久化和缓存
val rddCached = rddFromCollection.cache() // 默认将 RDD 缓存到内存
val rddPersisted = rddFromCollection.persist(StorageLevel.DISK_ONLY) // 将 RDD 持久化到磁盘
5、Spark 支持哪些语言的 API?每种语言的适用场景是什么?
scala:spark 底层实现的语言,用于开发 spark
java:企业级开发
python:利用 pyspark,结合 python 强大的数据分析库
r:sparkr api 支持统计分析和数据挖掘
sql:编写 etl
6、在 Spark 中,什么是 Transformation 和 Action?两者有什么区别?
Transformation:对数据进行转换,但不会立即执行计算,结果是一个新的 RDD,实际的计算只有在后续触发 Action 操作时才会执行。包括 map()、filter()、 flatMap()、groupBy()、reduceByKey() 等;
Action:触发实际的计算,并将结果返回给 Driver 或写入外部存储系统。会拉取之前定义的所有的 Transformation 并执行。包括 collect()、count()、saveAsTextFile()、take()、reduce() 等;
spark 的 Transformation 操作会记录到一个 DAG 中,直到遇到 Action 操作时才会真正执行,有助于优化调度和计算效率;spark 可以在执行前生成一个优化的物理执行计划,减少执行次数和数据传输;在进行多次 Action 操作时,若某个 Transformation 操作的结果需要被多次用到,可以用 persist() / cache() 持久化中间结果,减少重复计算的开销
val lines = sc.textFile("hdfs://data/sample.txt") // Transformation: textFile
val words = lines.flatMap(line => line.split(" ")) // Transformation: flatMap
val wordCounts = words.map(word => (word, 1)).reduceByKey(_ + _) // Transformation: map and reduceByKey
wordCounts.saveAsTextFile("hdfs://data/output") // Action: saveAsTextFile
7、Spark 的 DAG(有向无环图)是如何生成的?它在任务调度中的作用是什么?
DAG 是 spark 通过对用户编写的 高级 API 操作(如 map、filter、reduce 等)进行解析、优化和转化生成的,主要依赖 spark 的调度器和任务划分策略。
步骤:
a、用户在 Driver 节点上编写的代码经过逻辑闭包分析,生成 Stage
b、根据数据依赖关系将任务分解为每个 Stage,由窄依赖和宽依赖划分
c、每个 Stage 内部生成 TaskSet,这些 Task 并行执行
8、在 Spark 中,如何持久化 RDD?常见的持久化级别有哪些?
spark 中,一般通过 persist()、cache() 持久化 RDD,常见的持久化策略有:
a、MEMORY_ONLY:将 RDD 仅持久化到内存,如果内存不足会丢弃部分数据
b、MEMORY_AND_DISK:将 RDD 先存到内存,不足就存到磁盘
c、DISK_ONLY:仅持久化到磁盘
d、MEMORY_ONLY_SER:将 RDD 以序列化的形式持久化到内存,适用于数据量较大且计算不频繁的场景
e、MEMORY_AND_DISK_SER:将 RDD 以序列化的形式先存储到内存,内存不足就存到磁盘
持久化操作完成后,需要主动调用 unpersist() 方法释放缓存,以防占用资源过多
9、在 Spark 中,如何通过 cache() 和 persist() 优化性能?
通过 cache() persist() 缓存 RDD,避免多次重复计算,减少计算时间;对于机器学习算法等需要多次迭代的数据处理任务,缓存或持久化能显著提高计算效率
10、什么是 Spark 的惰性计算机制?它是如何工作的?
对数据集进行各种转换操作时,这些操作会被记录下来,但是不会立即执行,只有碰到动作操作时,才会执行所有之前的转换操作并生成最终结果。
优势:能够通过优化计划(任务合并、管道化操作等)提升性能,减少中间数据的生成,减少内存使用的 I/O 操作。
11、在 Spark 中,什么是分区?如何调整 RDD 的分区数量?
分区是数据处理(供并行处理)的最小工作单位。每一个分区可以在一个任务中被处理,调整 RDD 的分区数量主要有两种方法:
a、repartition(numPartitions: Int):增加或减少分区数量。但这会引起全局的 shuffle 操作,因此开销大
b、coalesce(numPartitions: Int, shuffle: Boolean = false):减少分区数量,可以选择是否进行 shuffle 操作,通过减少 shuffle 操作可以提高性能。
# 创建 RDD 时制定分区:在读取数据时,指定分区数量
val rdd = sc.textFile("hdfs://path/to/file", numPartitions)
# repartition 方法:重新调整分区数量,但会触发 shuffle,这意味着所有的数据会被重新分配到新的分区中
val repartitionedRDD = rdd.repartition(10)
# coalesce 方法:一般用来减少分区数量,可以避免 shuffle,加上 shuffle 参数是为了稀疏的分区也能合理合并
val coalescedRDD = rdd.coalesce(5)
分区数会影响并行度,资源使用,执行和调度开销
12、Spark 中的广播变量是什么?它在性能优化中的作用是什么?
是一种共享变量,允许在所有节点上高效地共享一个只读变量,而不是通过网络将该变量多次发送给各个节点,主要用于将任务中较小的数据集广播到各个节点,以避免反复传输。
# 广播变量的创建及使用:
val broadcastVar = sc.broadcast(Array(1, 2, 3))
val rdd = sc.parallelize(1 to 10)
val result = rdd.map(x => (x, broadcastVar.value.contains(x)))
广播变量的内部机制:广播变量会在 Driver 上进行序列化,并通过高效的广播算法将其分法到各个 Executor,每一个节点上存储一份拷贝,不仅减少网络传输,也减少了数据在节点间的传输距离。
常见使用场景:需要在多个任务中反复使用的大型只读数据集,比如参考数据表;相对较小且不会频繁改变的数据集,比如业务规则配置、字典等。
13、在 Spark 中,如何使用累加器来实现数据的聚合?
spark 中的累加器通过 Driver 创建,并在所有工作节点上对这个累加器变量进行操作,通过累加器可以对数据进行求和、计数等聚合操作。
# 创建累加器:在 Driver 中创建累加器
val sc = new SparkContext(conf)
val accum = sc.longAccumulator("AccumulatorName")
# 使用累加器:在操作 RDD 的各个 worker 中,通过累加器进行累加操作
rdd.foreach(x => accum.add(x))
# 获取累加器值:在操作完成后,通过 Driver 获取累加器的最终值
val result = accum.value
# 一个完成例子:
// 配置和创建 SparkContext
val conf = new SparkConf().setAppName("AccumulatorExample").setMaster("local")
val sc = new SparkContext(conf)
// 创建累加器
val accum = sc.longAccumulator("SumAccumulator")
// 创建一个 RDD
val rdd = sc.parallelize(Array(1, 2, 3, 4, 5))
// 在 RDD 的每个元素上使用累加器
rdd.foreach(x => accum.add(x))
// 获取累加器的值
println(s"Accumulated value is: ${accum.value}")
// 关闭 SparkContext
sc.stop()
累加器有 LongAccumulator、Double Accumulator、CollectionAccumulator
使用场景:计算综合、平均值、计数;记录调试信息或日志,能收集从所有 worker 节点传来的信息
限制:累加器的值只能在 driver 中读取,不能在 worker 中读;累加器是容错的,单个累加器的操作可能会应用多次,需要操作完成后对结果进行校验;累加器由 driver 中的变量管理,增加网络通信开销。
广播变量用于只读数据集共享,累加器用于共享可写变量。
14、Spark 如何与 Hadoop 的 HDFS 集成?它们之间的数据流动如何实现?
a、数据以文件的形式存储在 HDFS 的各个数据节点
b、当启动一个 spark 作业时,spark 通过 Hadoop 的 API 读取存储在 HDFS 上的数据,这个过程通过 Hadoop 的 InputFormat 类实现,它会根据文件格式将文件切分成若干分区,并为每个分区创建相应的输入分片
c、这些数据分片又会被分法给 spark executor,基于 spark 的 rdd 和数据分区来实现,每个分区的数据会被加载到 Executor 的内存中进行计算
d、当 spark 完成数据处理后,将结果写回 HDFS,通过 Hadoop 的 OutputFormat 类实现
15、Spark 中的 shuffle 操作是什么?它对性能有什么影响?
shuffle 是 spark 中的一种数据重分区过程,是在 stages 之间重新分配数据,使得一个阶段的输出成为下一个阶段的输入,涉及大量的数据序列化、网络传输和 磁盘 I/O; shuffle 分为两种:宽依赖和窄依赖,宽依赖需要跨节点传输数据,窄依赖只在一个节点内部进行,不涉及网络传输和 磁盘 I/O,典型的 shuffle 操作包括 reduceByKey、groupByKey、join 等
shuffle 的具体过程:各个分区处理数据并产生中间结果、中间结果被序列化并写入磁盘、中间结果通过网络传输到目标分区所在的节点、目标分区读取数据并进行反序列化,执行后续计算任务。
性能优化的技巧:
a、通过增加 task 并行度,减少单个 task 的数据量
b、避免宽依赖:尽量使用 类似 map、filter 等窄依赖操作
c、使用适当的操作:reduceByKey 比 groupByKey 更高效,因为它在 shuffle 之前先进行局部合并,减少了数据传输量
d、数据预聚合:使用 分区器 Partitioner 控制数据分布,提前进行数据预聚合,减少数据重组开销
e、优化选择高效的序列化方式如 Kryo
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("ShuffleExample").getOrCreate()
sc = spark.sparkContext
data = [("a", 1), ("b", 1), ("a", 1), ("b", 1), ("c", 1)]
rdd = sc.parallelize(data)
# 触发 shuffle 操作的 reduceByKey
result = rdd.reduceByKey(lambda x, y: x + y).collect()
print(result)
# 输出:[('a', 2), ('b', 2), ('c', 1)]
spark.stop()
16、在 Spark 中,如何避免 shuffle 操作?有哪些优化 shuffle 的方法?
避免 shuffle 操作:
a、使用 coalecse 代替 repartition:repartition 会触发完整的 shuffle,coalesce 通过减少分区数,有效避免 shuffle
b、mapPartitions 代替 map 和 filter:map 和 filter 操作时逐条处理数据,mapPartitions 允许批量处理数据,可以减少 shuffle 的发生,且高效利用内存
c、broadcast 变量:在 join 操作中,如果右表很小,可以使用 broadcast 变量将其广播到所有节点中,然后进行广播 join 以避免 shuffle
优化 shuffle 操作:
a、通过增加 task 并行度,调整 spark.default.parallelism 参数,减少单个 task 的数据量
b、优化选择高效的序列化方式如 Kryo
c、使用 DataFrame 或 Dataset API: 这二者对 shuffle 操作进行了优化,包括使用更紧凑的存储和执行计划优化
d、缓存和持久化:当一个 rdd 或 dataframe 被多次使用,可以使用 cache() / persist() 方法将它们缓存到内存中,避免重复计算
e、利用 speculations:Speculative Execution 可以以容错机制重新调度运行缓慢的任务,提升整体作业的性能,通过 spark.speculation 为 true 启用。
17、Spark 中的宽依赖和窄依赖是什么?它们有什么区别?
窄依赖:每个父 RDD 的分区最多对应一个 子 RDD 的分区
宽依赖:每个父 RDD 的一个分区对应多个 子 RDD 的分区,通常涉及到数据的重新分布,需要 shuffle
窄依赖因为没有 shuffle,所以数据不用跨节点传输,开销小;窄依赖故障恢复只涉及单个分区的数据,只需重新计算失败的分区数据即可,宽依赖则需要重新计算 shuffle 前的所有 父 RDD;
18、如何在 Spark 中优化 Join 操作?有哪些常见的优化策略?
广播小表:使用 BroadcastHashJoin ,减少通信消耗;
减少 shuffle 操作:通过对数据进行合适的分区和 coalesce 操作来减少 shuffle 带来的网络开销;
拆分复杂 join:将复杂的多表 join 操作拆分成多个单独的两表 join 操作,优化查询计划;
调整并行度;
使用合适的 join,比如 Sort-Merge Join 适用于大表和大表的 join,通过对表进行排序和合并;Broadcast Hash Join 适用于小表和 大表的 join 操作,将小表广播到所有工作节点,避免 shuffle;Shuffled Hash Join / Shuffle Sort Merge Join 适用于大规模需要 shuffle 的场景;
基于业务逻辑定义合适的 partitioner,使用 repartition 或 coalesce 操作调整分区数量;
使用合适的存储格式如 Parquet、ORC 可以提高数据读取和写入性能;
预先做一次数据梳理操作,将数据分割为更加适合 join 的格式;
将频繁使用的表或 join 操作的中间结果缓存到内存,避免重复计算;
调优参数:spark.sql.shuffle.partitions 增加 partition 数量;spark.executor.memory 配置更多的内存
19、在 Spark 中,如何使用 repartition 和 coalesce 进行分区调整?它们有什么区别?
repartition 可以增加或减少分区,但会触发 shuffle;coalesce 通常用于减少分区,可以省去 shuffle;
repartition 一般用于增大并行度的计算任务,比如 CPU 密集型任务,或对数据进行全局重新排序时用到;coalesce 通常用于 I/O 密集型任务中,减少任务调度和数据读取的开销
20、Spark 的任务调度机制是如何工作的?如何根据集群的资源情况进行任务调度?
任务调度机制主要负责将作业拆分成多个任务并分配给集群中的各个工作节点去执行,依赖集群管理器,如 yarn、mesos 等等来完成资源申请和调度,实现高效的资源利用和任务并行化执行;
spark 任务调度的过程分为以下几个步骤:
a、作业提交:用户通过 sparkcontext 提交作业
b、阶段划分:DAG Scheduler 将作业解析为一个由 RDD 转换而来的 DAG,然后将 DAG 拆解为多个 Stage 和 Task
c、任务分配:Task Scheduler 接收阶段信息,将任务分配到具体的 Executor
d、任务执行:Executor 根据接收的任务和数据分片进行计算
e、结果反馈:任务执行完成后,结果返回给 Driver
默认情况下,spark 使用 fifo 调度策略,可以使多个 作业公平竞争资源