全文共 4,000 字 预计阅读 12 分钟
bg

spark20问-3

1、Spark 的 Shuffle 写阶段是如何工作的?如何优化 Shuffle 的写入性能?

shuffle 阶段涉及大量数据的分发、排序和合并,shuffle 写阶段最主要的任务就是将上游任务的输出数据分发到不同的 reducer,以便后续处理;

a、executor 执行的每个 task 根据 partitioner 将数据分区

b、将分区的数据序列化成字节流,以便后续的网络传输和磁盘 I/O 操作

c、将序列化后的数据存储在本地磁盘的临时文件中,为每个分区创建一个临时文件

d、任务完成时,会向 Driver 发送包含每个块的位置和大小的信息,以便下游任务能够正确地读取这些数据

优化 shuffle 的写入性能:

a、使用 kryo 序列化

b、合理调整 rdd 的分区数,也就是 task 的数量,避免过度分区或者分区过少,分区过多会导致文件过多,频繁 i/o,分区过少造成单个分区数据量过大,产生瓶颈

c、调整 spark 的配置参数,spark.shuffle.file.buffer (默认 32KB,可以增大以减少写入磁盘的频次)和 spark.reducer.maxSizeInFlight(默认 48MB,控制 shuffle 数据在网络传输线程间的内存占用,适当提高可以提升网络传输的效率) 以优化写入缓冲区大小,减少磁盘 I/O 的次数

d、对于小文件较多的场景,启用 spark 的合并文件特性,减少文件数量,提升文件写入效率与读取效率

2、Spark 的 Stage 划分机制是如何设计的?如何优化 Stage 的划分以提升任务执行效率?

stage 划分机制为依赖关系驱动,每次有一个 action 会启动一个 spark job,该 job 被划分为多个 stage,每个 stage 由一系列无依赖关系的任务组成,可并行执行;stage 通常在窄依赖中划分,当出现宽依赖才会新建 stage;

优化 stage 可以从:减少 shuffle、调整并行度、广播变量、缓存中间结果

3、Spark 的资源调度器是如何工作的?如何调优资源调度策略?

资源调度器负责将应用程序的任务合理地分配到各个节点上执行,主要工作机制包括:

a、Driver 向 集群管理器 yarn、kube 申请资源

b、集群管理器分配资源后,Driver 根据 DAG 确定的优先级将任务分配到各个 executor 上

c、executor 在接收到任务之后,按照任务的依赖顺序来执行具体的计算工作

d、Driver 负责监控任务的执行情况,出现失败会重新分配和调度任务

为了调优资源调度策略,可以从:

a、合理配置 executor 的数量和大小:调整 spark 配置文件中的 spark.executor.instances 和 spark.executor.memory 等参数

b、开启动态资源分配 spark.dynamicAllocation.enabled,根据任务需求动态调整资源

c、调整任务并行度 spark.default.parallelism 和数据分区的数量,提高资源利用率

d、配置公平调度器 spark.scheduler.mode=FAIR,在多任务环境下更好地分配资源

4、Spark 的任务重试机制是如何实现的?如何通过任务重试提高容错能力?

a、当一个任务执行失败时,spark 会捕获该任务失败的异常信息

b、spark 判断该任务是否可以被重试

c、如果可以重试,spark 为该任务分配新的计算资源并重新执行

d、spark 的任务重试机制有一个预设的最大重试次数,默认 4 ,通过 spark.task.maxFailures 参数配置,超过这个次数任务仍然失败的则判定为无法恢复, spark.stage.maxConsecutiveAttempts 核定一个 stage 可以重试的次数,默认也是 4

5、在 Spark 中,如何通过自定义 Partition 实现数据分区优化?

# 自定义 Partition 类,负责分区的索引管理:
import org.apache.spark.Partition

class CustomPartition(val index: Int) extends Partition {
  override def hashCode(): Int = index
}

# 自定义分区器类,其中 numPartitions 方法返回分区数,getPartitions 方法根据键的哈希值将其分配给一个特定的分区:
import org.apache.spark.Partitioner

class CustomPartitioner(partitions: Int) extends Partitioner {
  override def numPartitions: Int = partitions

  override def getPartition(key: Any): Int = {
    // 定义你的分区逻辑,这里简单地使用了 key 的哈希值
    key.hashCode % numPartitions
  }
}

# 应用自定义分区器:记那个自定义的分区器应用到 rdd 上
import org.apache.spark.SparkContext
import org.apache.spark.SparkConf

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

val rdd = sc.parallelize(Seq((1, "a"), (2, "b"), (3, "c"), (4, "d")), 4)
val partitionedRDD = rdd.partitionBy(new CustomPartitioner(2))

// 查看每个分区的数据
partitionedRDD.mapPartitionsWithIndex((index, iterator) => iterator.map((index, _))).collect.foreach(println)

6、Spark 中的 Tungsten 引擎是什么?它如何通过物理执行优化提升性能?

是 spark 中的物理执行引擎,目的是通过低级别的物理执行优化来提升计算性能,主要关注的是 cpu 和内存的高效使用,

a、通过更智能的内存管理机制,如使用二进制存储格式、分层内存池等,减少了 jvm 垃圾回收的影响,提升了性能

b、通过向量化的计算方式,即在处理数据时批量操作多行数据,减少了 cpu 的开销,减少指令数量,提升了执行效率

c、通过运行时代码生成技术将 spark sql 任务编译为高效的字节码,使得执行计划更高效,充分利用现代 cpu 的性能特性

7、在 Spark 中,如何优化大规模数据集上的 Join 操作?有哪些实际应用场景?

a、当一个数据集较小,另一个数据集较大时,可以采用广播 join 的方式将较小的数据集广播到所有执行节点,本地执行 join 操作,避免了大型 shuffle 操作

b、在执行 join 操作前,确保数据集按照相同的键进行分区,可以减少数据传输和 shuffle 的开销

c、若数据集由大量小文件组成,先进行合并操作,大量存在会导致大量的小任务,增加调度开销和 IO 成本,从而拖慢 Join 操作

d、对于多个操作使用相同的中间结果,可以将中间结果进行缓存,然后再多次使用

e、在 join 操作中调优 spark.sql.shuffle.partitions 参数,利用合适的 shuffle 分区数有效平衡任务负载,避免某些数据分区过于集中

val products = spark.read.parquet("hdfs://path_to_products")
val logs = spark.read.parquet("hdfs://path_to_logs")
val broadcastProducts = broadcast(products)
val joinedData = logs.join(broadcastProducts, "productId")

val rawData = spark.read.json("hdfs://path_to_raw_data")
val consolidatedData = rawData.coalesce(100)  // 合并成较少的分区
consolidatedData.write.parquet("hdfs://path_to_cleaned_data")

val sales = spark.read.parquet("hdfs://path_to_sales")
val customers = spark.read.parquet("hdfs://path_to_customers")
val salesWithRegion = sales.join(customers, "customerId").cache()

val regionReport = salesWithRegion.groupBy("region").sum("salesAmount")
val productReport = salesWithRegion.groupBy("productId").sum("salesAmount")

8、Spark 的内存管理分为哪些部分?如何优化内存管理以提高作业性能?

分为 执行内存和存储内存,执行内存用于存放中间结果和 task 级的数据结构,存储内存用于缓存 rdd 和 共享变量,优化内存管理,可以通过:

a、调整 spark.executor.memory 参数,适当增加 executor 内存的大小,确保有足够内存来处理任务

b、调整 spark.memory.fraction 参数,控制用于存储和执行的内存比例,通过调节它可以优化内存的分配

c、根据数据的重要性和使用频率选择合适的持久化级别,防止不必要的高内存占用

d、使用 unpersist() 方法释放不再需要的缓存数据以腾出内存空间

e、设置适当的并行度参数 spark.default.parallelism,防止单个 task 占用过多内存

spark 的内存管理采用了统一的内存管理模型,存储内存和执行内存可以互相借用;

运行大规模任务,为了防止出现内存溢出的情况,可以适当调整 spark.memory.storageFraction 和 spark.memory.executionFraction 参数来调节二者的比例;

9、在 Spark 中,如何通过调整并行度和任务划分来优化执行效率?

调整并行度和任务划分是为了让资源充分利用、任务平衡分布、减少数据倾斜,基础是使用 spark.sql.shuffle.partitions、repartition 和 coalesce

a、spark.sql.shuffle.partitions 决定 shuffle 操作后的 分区数,默认为 200,可以根据数据规模和集群资源调整

b、当数据分布不均衡时,可以用 repartition(n) 将数据重新分为 n 个,或者用 coalesce(n) 来减少分区数。

c、避免数据倾斜:hash partitioner、salting、aoe、适当增大并行度,增加分区数来分散数据压力

d、在调整并行度和任务划分时,要注意内存和资源的平衡

e、使用 spark web ui 工具,实时监控集群资源利用率

10、Spark 的动态分区调度是如何工作的?它对资源利用率有什么影响?

动态分区调度允许任务动态调整分区数量,可根据实际运行时的资源情况和负载情况,实时调整。

a、通过监视每个分区的处理时间,将负载重新分配给速度较快或负载较轻的节点,防止某些节点成为瓶颈

b、根据资源的使用情况进行调整,提高资源的利用率

c、缓解数据倾斜的问题

11、Spark 的 DAG 执行计划是如何生成的?如何优化 DAG 以减少执行开销?

a、logical plan:用户编写的代码被转换为一个逻辑计划,是数据集操作的高级描述

b、optimized logical plan:这个逻辑计划经过 catalyst 优化器的优化,生成一个优化的逻辑计划,包括谓词下推,列剪裁等优化

c、physical plan:优化后的逻辑计划接着会被转化为一个或多个物理计划,物理计划描述了具体的执行步骤,包括生成的 rdds 以及如何针对这些 rdds 进行计算

d、spark 通过代价模型选择一个代价最低的物理计划来执行

常见优化:减少 shuffle、缓存中间结果、提升并行度、数据倾斜处理、减少宽依赖

catalyst 是 spark sql 的查询优化器,关注于通过多种规则和代价模型来优化查询计划,能自动进行多种优化操作,例如谓词下推、合并过滤条件等,这使得查询速度更快、执行效率更高

spark 采用 job、stage、task 三级模型来进行分布式计算,一个作业被拆为多个阶段,每个阶段有被进一步分为多个任务,每个任务将在一个节点上进行

数据局部性:尽量使任务在数据所在的节点上执行

谓词下推指将过滤条件下推到数据源层面,减少需要扫描和处理的数据量,在 sparksql 中,通过 dataframe api 附带的过滤条件可以自动被下推

在实际的调优中,可以利用 spark 提供的 web ui 来监控和分析作业的执行情况,找到瓶颈和潜在优化点

12、在 Spark 中,如何处理 Executor 和 Driver 之间的通信瓶颈?

通信瓶颈通常是大量数据传输或者频繁的通信请求引起的

a、尽可能减少传输的数据量,例如通过在 executor 端进行数据预处理,避免将大量中间数据传回 Driver

b、增加 Driver 的内存和 cpu 配置,以应对更高的通信负荷

c、使用广播变量

d、减少 shuffle 操作

e、使用缓存

f、优化序列化

g、数据局部性优化:通过确保数据处理尽量在数据本地或附近的节点进行

13、Spark 中的 Windowing 操作是如何实现的?它的应用场景有哪些?

通过 dataframe api 的 window 函数实现,windowing 操作允许用户通过指定的窗口的范围对数据进行分组、排序并聚合,应用场景包括:时间序列分析、滑动窗口计算、频率计算、排名

spark 中窗口函数分为三类:聚合函数、排名函数、分析函数

# 窗口的定义
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

val windowSpec = Window.partitionBy("partitionColumn").orderBy("orderColumn").rowsBetween(Window.unboundedPreceding, Window.currentRow)

# 有一个包含用户点击的 dataframe,要统计每个用户过去 10 条点击记录的总点击数
val clicks = spark.read.json("/path/to/clicks.json")

val windowSpec = Window.partitionBy("userId").orderBy("timestamp").rowsBetween(-10, 0)

val clickCounts = clicks.withColumn("total_clicks", sum("clickCount").over(windowSpec))

14、Spark 的内存和磁盘溢写策略是如何设计的?如何优化以避免频繁的溢写?

内存和磁盘溢写主要用于处理内存不足,将数据溢写到磁盘的情况。通过 spark 的存储级别和内存管理策略来实现

a、spark 提供了多个存储级别来决定数据的持久化方式

b、spark 的内存管理分为 execution 内存和 storage 内存,分别用于执行操作和缓存数据

15、Spark 的 Shuffle 读取阶段是如何优化的?如何减少网络 IO 和延迟?

a、通过数据合并减少文件的数量和大小,降低网络传输的次数

b、尽量利用已存储在本地的中间结果,减少跨节点的数据传输

c、将数据的读取和写入操作串联起来,形成管道,使得数据不需要等待所有数据都读取完毕后才开始写入

d、一次性批量读取较大的数据块,而不是小块多次读取,减少网络请求的次数和相关的开销

e、优化 shuffle 文件,比如文件缓冲及溢出、索引文件

f、采用高效的数据压缩技术

g、利用 netty 高效进行网络传输,shuffle 使用 异步 io 来处理大量的网络请求,提升传输的速度

16、在 Spark 中,如何通过动态资源分配实现资源的精细化管理?

spark 根据实际的作业需求增加或释放 executor,从而更有效地利用集群资源并控制成本和延迟

a、启用动态资源分配功能:spark.dynamicAllocation.enbaled 为 true

b、配置最小和最大 executor 数量:spark.dynamicAllocation.minExecutors 和 spark.dynamicAllocation.maxExecutors 定义 executor 数量的上下界

c、设置空闲超时时间 spark.dynamicAllocation.executorIdleTimeout 来配置 Executor 在空闲状态下等待被释放的最大时间

d、启用 External Shuffle Service:在 Yarn / Mesos 等资源配置管理工具中,配置 External Shuffle Service 来保存中间数据,即使 executor 被释放 ,shuffle 文件依然可用

17、Spark 的 RDD 转换为 DataFrame 时有哪些性能优化策略?

a、从 rdd 转换到 dataframe 时,可以手动定义 schema 而不是依赖 spark 的自动推断类型,减少 spark 在推断类型上的开销

from pyspark.sql.types import StructType, StructField, IntegerType, StringType

schema = StructType([
    StructField("id", IntegerType(), nullable=True),
    StructField("name", StringType(), nullable=True)
])

rdd = sc.parallelize([(1, "Alice"), (2, "Bob")])
df = spark.createDataFrame(rdd, schema)

b、尽量避免在 rdd 和 dataframe 之间的反复转换,这种操作会反复执行序列化和反序列化,增加计算开销

c、在转换之前,利用 cache() 或 persist() 方法将原始 rdd 或 dataframe 缓存到内存中,以便多次使用时避免重复计算

d、尽量避免在转换过程中引入宽依赖

18、Spark 的容错机制是如何设计的?它在大规模数据处理中的作用是什么?

容错机制主要依赖于两大核心概念:rdd 和 dag,每个 rdd 都被视为一个有容错特性的集合,当 rdd 一部分实效,spark 可以通过 lineage 重新计算丢失的数据。spark 通过标准的重试机制和任务推测执行机制来处理这些问题。

rdd 的血统:记录了 rdd 是如何从其他 rdd 派生而来的信息

dag:spark 会将 job 划分为 stage,每个stage 由一系列 task 组成,每个 rdd 的计算会形成一个 dag,通过这个 dag 可以知道每个数据的依赖关系

19、Spark 中的 Checkpoint 机制如何实现数据恢复和任务重启?它对性能有什么影响?

checkpoint 通过把 rdd 的 lineage graph 进行阻断,把当前的 rdd 数据写入分布式存储系统来保存快照

任务重启时,spark 会从 checkpoint 文件中读取数据,而不必重新计算被 checkpoint 的 rdd

对性能的影响主要体现在:

a、写入 checkpoint 文件的过程中,需要将数据从内存 / 磁盘传输到分布式存储系统,导致额外的 IO 开销和网络传输

b、会引发一些延迟

20、Spark Structured Streaming 如何保证 Exactly Once 语义?它的底层实现是什么?

a、幂等写入:多次执行同一个操作结果一致,spark 利用目标存储系统的特性,比如 kafka 幂等生产者 api,确保消息写入即使重试也不重复

b、事务性存储:利用事务性存储系统,如 kafka delta lake

c、检查点机制:通过 checkpoint 记录流处理中的元数据信息,包括操作的偏移量、状态等,将这些信息保存在可靠的存储系统中

Back to Blog