全文共 4,235 字 预计阅读 13 分钟
8gu

spark20问-2

1、在 Spark 中,如何通过调整并行度来提升性能?

调整并行度来提升性能主要通过调整 RDD 和 DataFrame 的分区数量来实现。

调整并行度具体有以下 3 种方式:

a、配置默认并行度

# 在建立 spark context 时,通过配置参数来调整默认并行度
val conf = new SparkConf().setAppName("MyApp").setMaster("local[4]")
conf.set("spark.default.parallelism", "12")
val sc = new SparkContext(conf)
local[4] 表示在本地运行,用 4 个线程模拟 4 个核心,
spark.default.parallelism 被设置为 12,默认并行度就为 12

# 调整 RDD 的分区数量:
val rdd = sc.textFile("hdfs://path/to/file", 12) // 12 分区
val repartitionedRDD = rdd.repartition(20) // 将 RDD 重分区为 20
val coalescedRDD = rdd.coalesce(6) // 将 RDD 合并成 6 个分区
repartition 会触发 全量 shuffle,而 coalesce 只会合并分区

# 调整 DataFrame 的分区数量:
val df = spark.read.json("hdfs://path/to/json")
val repartitionedDF = df.repartition(12) // DataFrame 分区数
val coalescedDF = df.coalesce(6) // DataFrame 合并分区

此外,还需要注意:

a、数据倾斜问题:调整并行度的同时要注意避免某些分区的数据过多 / 过少的情况

b、要根据资源调整并行度,确保任务能够在合理的时间内完成,不能因为资源的争抢带来性能瓶颈

2、Spark 的内存管理机制是如何设计的?如何优化内存的使用?

spark 的内存管理机制分为静态内存管理和 统一内存管理;为了优化内存的使用,减少内存溢出和 gc 的开销,需要注意参数调优和内存分配的策略;

a、缓存区:用于缓存 RDD、DataFrame 和 Dataset 的数据,以便于加速后续的迭代计算;

b、执行内存:用于存储在计算过程中生成的中间数据,比如 shuffle、join 操作需要用到的内存。

在统一内存管理中,缓存区和执行内存共同使用一个内存池,内存可以动态分配。

优化内存使用:

a、调整配置参数:spark.executor.memory、spark.memory.fraction、spark.memory.storageFraction 等,比如执行大型操作时,如 shuffle,可能导致内存溢出,可以增加 spark.shuffle.memoryFraction 或增大 spark.executor.memory 参数,并合理划分数据分区

b、使用缓存:适当缓存和持久化关键的 RDD,减少重复计算

c、使用高效的序列化库,如 Kryo,减少数据在内存中的占用

d、gc 调优:通过配置参数(-XX:NewRatio、-XX:SurvivorRatio、-XX:MaxPermSize 等),或使用 不同的 gc 算法

e、使用 spark ui 或其他监控工具,实时监控内存使用情况

f、当 RDD 或 DataFrame 持久化时,选择不同的存储级别

g、合理设置任务的并发度以及执行顺序,避免过多的任务同时争抢内存资源,通过 spark.sql.shuffle.partitions

3、在 Spark 中,如何监控作业的执行?有哪些常用的监控工具?

a、使用 spark ui,在 web 界面上查看作业的执行情况、任务分布、失败节点等信息

b、如果 spark 作业运行在 yarn 上,可以使用 yarn 的 RM Web UI 来查看资源使用和作业状态

c、Ganglia 和 Graphite 可以收集展示更多的系统指标,比如节点 CPU、内存使用等详细信息

d、通过 spark 的日志系统,可以查看每个阶段的详细日志信息,一般包括 INFO、WARN、ERROR 等详细日志

e、spark history server:保存和展示已经完成的 spark 作业的执行历史记录

4、Spark 中的广播变量和累加器有什么区别?它们在不同场景中的应用分别是什么?

广播变量用于将集中的只读数据从 Driver 有效分发到每个节点,因为不是每次任务发送时都序列化和发送,大大减少了数据传输的开销,提高了性能;

累加器用于执行 spark 中统计、计数等聚合操作。只能在 Executor 上增加,而 Driver 指可以读取它,不能修改,避免了并发修改的问题。

5、在 Spark 中,如何使用窗口操作处理实时数据流?

可以利用 spark streaming 和 structured streaming 对实时数据流进行窗口操作。

a、定义窗口大小和滑动间隔:窗口大小决定了每个窗口包含的数据时间范围,滑动间隔决定了窗口的移动频率。

b、应用窗口操作:使用 windows 函数或方法对数据进行窗口化处理

c、执行聚合或其他操作:在窗口划分完成后,可以对窗口内的数据进行各种聚类操作,比如求和、取平均值

from pyspark.sql import SparkSession
from pyspark.sql.functions import window

spark = SparkSession.builder.appName("WindowedStreaming").getOrCreate()

# 创建一个 Streaming DataFrame 从一个数据源,比如 Kafka
lines = spark.readStream.format("kafka").option("subscribe", "topic").load()

# 假设 lines DataFrame 有一个时间戳列 "timestamp" 和一个值列 "value"
# 先把数据流转换成标准的数据格式
lines = lines.selectExpr("CAST(value AS STRING)", "timestamp")

# 定义窗口大小和滑动间隔
windowedCounts = lines.withWatermark("timestamp", "10 minutes") \
                      .groupBy(window("timestamp", "10 seconds", "5 seconds")) \
                      .count()

# 开始查询,并将结果输出到控制台
query = windowedCounts.writeStream.outputMode("update").format("console").start()

query.awaitTermination()

watermark:数据如果乱序到达,可以设置一个水位线表示延迟容忍的最大时间,超过这个时间的数据将被丢弃

窗口类型:翻滚窗口(大小固定、不会重叠)、滑动窗口(大小固定,部分重叠)、会话窗口(根据数据的活动间隔动态生成窗口)

可以合理设置 spark.streaming.batch.duration 以平衡延迟和处理效率

6、Spark 的 Checkpoint 机制是什么?它在大规模数据处理中的作用是什么?

checkpoint 机制是指将 rdd 或者 dataframe 的中间结果持久化到稳定存储中,以便在出现故障时可以从 checkpoint 恢复计算,而不是从头计算,这在大规模数据处理中尤为重要;

在长时间运行的 spark 作业中,使用 checkpoint 可以防止由于节点故障导致的长时间计算重启;

在调用 rdd 的 checkpoint 方法时,spark 会对该 rdd 进行物化,并将结果保存到指定的存储位置,在 streaming 场景中,可以设置周期性的 checkpoint

7、Spark 的 DAG Scheduler 和 Task Scheduler 分别是什么?它们的作用是什么?

dag scheduler 时 spark 中的调度器,将用户提交的应用程序划分为多个 stage,每个 stage 由一组 task 组成,并依赖于其他 stage 的结果。dag scheduler 负责生成 dag,划分和分配计算任务至不同的 stage(每一个 stage 都是一个 shuffle 的边界),并处理 stage 之间的依赖关系;

task scheduler 是 spark 中另一个调度器,接收 dag scheduler 划分的任务之后,将这些任务分配到集群中的各个节点来执行,task scheduler 负责确保任务的高效调度、资源管理、故障恢复等。

dag scheduler 更关注全局任务划分,task scheduler 关注具体任务的执行和资源分配;负责处理任务之间的依赖关系并确保这些依赖关系在执行过程中得到满足,而 task 关注具体任务执行中的依赖解析和管理

8、在 Spark 中,如何优化数据的序列化和反序列化过程?有哪些常用的序列化方法?

使用 kryo 序列化库,比默认的 java 序列化库更高效;

kryo 运行注册自定义类来提升序列化性能,通过 registerKryoClasses 方法可以注册需要的类:

val registeredClasses = Array(
   classOf[YourClass1],
   classOf[YourClass2]
)
conf.registerKryoClasses(registeredClasses)

spark.kryo.registerationRequired:boolean 是否要求注册 Kryo 类,true 的话未注册的泪将导致作业失败

spark.kryo.unsafe:boolean 是否启用 Kryo 的不安全模式,可以提高 Kryo 的性能,但会牺牲一定的安全性

spark.kryoserializer.buffer.max:设置序列化缓冲区的最大允许大小

9、Spark 中的 Fault Tolerance 机制是如何设计的?如何保证任务的容错性?

Fault Tolerance 机制主要通过 数据的容错和 任务的容错来具体实现;

数据的容错是通过 rdd 来实现的,每个 rdd 包含一些列的分区以及这些分区的计算方式。当某个分区的数据丢失时,可以根据这些依赖来重算丢失的数据,rdd 的容错性时通过 lineage 机制来保证的,它记录了 rdd 从原始数据经过一系列转换操作得到当前数据集的计算过程。

任务的容错时通过重新执行失败的 task 来实现的,每个 job 被划分成多个 task,每个 task 处理一个 rdd 分区,当一个任务执行失败或所在的节点宕机的时候,spark 会重新提交该任务到其他可用节点执行,spark 也支持 speculative execution 机制,在某些任务进展缓慢时,尝试在其他节点上启动该任务的备份以竞赛完成;

其他的容错机制:hdfs 的多副本存储;spark 的错误重试机制,最大次数可以通过 spark.task.maxFailures 配置;checkpoint 可以将 rdd 持久化到 磁盘,即使整个应用程序失败了,也可以从 checkpoint 恢复 rdd,不需要从头计算所有的数据。

10、在 Spark 中,如何通过调整数据分区数提高作业执行效率?

可以通过 repartition(numPartitions) 对所有数据进行重新分区,用于增加数据分区数,可以更加均匀地分配数据,但会引入较大的 shuffle 开销;

coalesce(numPartitions) 这个方法通常用于减少数据分区数量,不会触发 shuffle,因此更高效

调整分区数并不是万能的,需要结合实际场景和数据特点做出最优选择:

a、分区数要与集群资源匹配,尽量匹配集群的核心数,比如集群有 40 个 cpu 核心,分区数可以设置为 2-3 倍左右

b、避免数据倾斜:如果某些分区的数据量明显偏多,会导致数据倾斜问题,可以通过增加分区数和使用自定义 partitioner 来缓解

c、jacobian shuffle:在可能的情况下减少 shuffle 的数据量

d、结合广播变量、cache / persist 方法、使用适当的存储级别等等

例如在实际开发中,遇到处理一个 2TB 日志文件的任务,可以合理设置分区数、使用 persist(StorageLevel.MEMORY_AND_DISK) 缓存中间结果、尽量避免不必要的 shuffle、结合 broadcast 变量优化小表广播

11、在 Spark 中,如何处理数据倾斜问题?有哪些常见的优化策略?

a、增加随机前缀或后缀:适用于大量数据集中在少数键值时,可以在 map 阶段给 key 添加随机前缀,将这些数据打散到不同的 reducer 处理,最后在 reduce 阶段去除前缀,恢复原始 key

b、广播小表:适用于两个表进行 join 操作时,一个表非常小,一个非常大。通过 broadcast 操作把小表数据分发到所有节点,这样每个节点都可以独立完成 join 操作

c、salting:适用于数据处理中出现高度倾斜的数据,比如点击流日志

d、分区策略调整:适用于数据倾斜由于默认分区器没有均匀分布数据的情况,通过自定义分区器,让其来实现 key 的均匀分布,比如 hashpartitioner

e、增大并行度:适用于倾斜主要表现在任务量较少的数据,通过调大 spark.sql.shuffle.partitions 来增加并行处理的分片数

f、数据预处理:适用于源数据本身存在问题,比如大量小文件或者存在明显倾斜,通过对数据预聚合、过滤、转化来实现

12、在 Spark 中,如何通过 DAG 调度优化任务执行?有哪些具体优化策略?

a、减少宽依赖的数量,增加窄依赖

b、优化 stage 的划分,减少 stage 的数量

c、缓存和持久化中间结果,避免重复计算

d、合并和重用 rdd 的转换操作,减少生成中间 rdd 的开销

e、避免数据倾斜,均衡任务负载

13、在 Spark 中,如何使用 GraphX 进行图计算?GraphX 的应用场景有哪些?

# 引入相关的库和初始化 spark 环境
import org.apache.spark._
import org.apache.spark.graphx._
import org.apache.spark.rdd.RDD

val conf = new SparkConf().setAppName("GraphXExample")
val sc = new SparkContext(conf)

# 创建顶点和边:
val vertexArray = Array(
  (1L, ("Alice", 28)),
  (2L, ("Bob", 27)),
  (3L, ("Charlie", 65)),
  (4L, ("David", 42)),
  (5L, ("Ed", 55))
)

val edgeArray = Array(
  Edge(2L, 1L, 7),
  Edge(2L, 4L, 2),
  Edge(3L, 2L, 4),
  Edge(3L, 4L, 3),
  Edge(4L, 1L, 1),
  Edge(5L, 3L, 6)
)

val vertexRDD: RDD[(Long, (String, Int))] = sc.parallelize(vertexArray)
val edgeRDD: RDD[Edge[Int]] = sc.parallelize(edgeArray)

# 根据顶点和边创建图:
val graph = Graph(vertexRDD, edgeRDD)

# 执行图计算操作,比如计算每个顶点的度数:
val degrees = graph.degrees.collect
degrees.foreach { case (id, degree) => println(s"Vertex $id has degree $degree") }

应用场景包括:社交网络分析、交通网络、推荐系统

14、在 Spark 中,如何通过动态资源分配优化集群的资源使用效率?

开启动态资源分配之后,spark 会根据作业的需求动态调整执行器的数量,不仅能提高资源利用率,还能根据作业负载的变化及时增加或减少资源

a、启用动态资源分配:通过 配置 spark.dynamicAllocation.enabled 为 true,开启了 spark 的自动资源配置功能

b、设置最小和最大执行器数:通过 spark.dynamicAllocation.minExecutors 和 spark.dynamicAllocation.maxExecutors 配置

c、配置空闲时间:设置 spark.dynamicAllocation.executorIdleTimeout 控制空闲执行器的回收时间,避免资源浪费,比如 60s

d、资源管理器的配置:确保 yarn、mesos 支持并启用了动态资源分配

15、在 Spark 中,如何实现异步操作?异步操作对性能优化有什么帮助?

通过 future 和 promises 来实现,避免了阻塞作业,提高了程序的并行度和吞吐量。

其中,future 表示一个未来可能完成或失败的操作,可以在另一个线程计算值时,立即返回一个代表该计算结果的对象; promises 是由开发者手写的 futures,提供一个方法来填充 future 的值或状态

可以一步读取或写入一个数据源

import scala.concurrent.Future
import scala.concurrent.ExecutionContext.Implicits.global

def longRunningTask(): Int = {
  Thread.sleep(10000)
  42
}

val futureResult = Future {
  longRunningTask()
}

futureResult.onComplete {
  case Success(value) => println(s"Task completed successfully with result: $value")
  case Failure(e) => println(s"Task failed with exception: $e")
}

16、Spark 中的推测执行机制是什么?它在任务执行中起到什么作用?

当 spark 检测到某个 task 执行时间明显比其他 task 要长时,就会启动一个备用 task,并行地重试执行这个过长的 task,以加速整体任务的完成;在任务执行中,推测执行机制可以显著降低因少数执行时间过长的 task 导致整个 job 延迟完成的情况,通过启动备用任务,如果新任务比原先的任务更快完成,那么原先的任务就会被取消;

spark.speculation.interval:检测任务的时间间隔,默认时 100ms

spark.speculation.multiplier:用于决定一个任务是否需要推测执行的阈值系数,默认为 1.5,这意味着任务执行时间必须是其他任务执行时间的 1.5 倍以上才会被推测执行

spark.speculation.quantile:指定推测执行任务的比例,执行时间最慢的前几百分位任务会被推测执行

17、在 Spark 中,如何优化内存管理和数据溢写问题?

a、调整内存参数:spark.executor.memory spark.driver.memory 来控制执行器和 驱动程序的内存分配

b、利用 cache 和 persist 机制,在执行反复计算的操作时,将数据持久化到内存或磁盘

c、选择高效的序列化方式

d、合理设置数据的分区数量和大小,通过 repartition 和 coalesce 方法调整数据分区

e、调整 jvm 的垃圾回收参数 spark.executor.extraJavaOptions,优化 GC 停顿时间

f、当内存不足以存储所有数据时,可以控制数据溢写到磁盘,spark.storage.memoryFraction 和 spark.shuffle.memoryFraction 来优化内存与磁盘的比例

18、在 Spark 中,如何利用广播变量优化 Join 操作?它的性能提升原理是什么?

将小的数据集使用 sparkContext.broadcast() 方法广播,在进行 join 操作时,用广播变量而不是直接使用原始的 rdd 或 dataframe

在传统的 join 中,两个数据集都需要在各个节点之间进行 shuffle 操作,导致了大量的数据传输和排序工作,开销较大,广播变量较小的数据集只需传播一次,随后可直接在本地使用,大大减少了网络传输和 shuffle 操作

19、Spark 中的 Structured Streaming 是什么?它与 Spark Streaming 有什么区别?

依赖 spark sql 进行流式数据处理,使用 dataframe 和 dataset api,允许用户以声明式的方式来处理流式数据;spark streaming 相比而言更旧,区别在于 后者基于 dstream api;前者提供更高层次的抽象和端到端的一致性保证,spark streaming 更接近于底层操作;前者有更好的容错机制,可以提供精确一次的语义保证,后者只能至少一次;前者有更好的性能和更低的延迟,因为它利用了 sparksql 的优化和查询计划

20、Spark 中的 Catalyst 优化器是如何基于代价模型优化查询计划的?

catalyst 优化器通过一系列的规则和代价估算,对用户提交的查询计划进行语法解析、逻辑计划优化、物理计划选择,然后生成高效的执行计划:

a、解析:将 sql 语句解析为 抽象语法树 ast

b、分析:对 ast 进行语义分析,包括列的解析、表的解析、类型检查,生成未优化的逻辑计划

c、逻辑优化:应用一系列逻辑优化规则来简化和优化未优化的逻辑计划,规则包括但不限于谓词下推、投影下推、列剪裁、子查询去相关化

d、物理规划:为每个逻辑计划生成多个可能的物理执行计划,每个物理计划表示 sql 查询实际运行时的执行方式,包含具体的执行操作,比如扫描操作和 join 操作。

e、代价评估:使用代价模型评估每个物理执行计划的代价,选择代价最小的方案作为最终的执行计划,考虑因素包括 I/O 成本、CPU 成本、网络开销等

Back to Blog