Spark 核心

共 19 题
📑 题目列表 19 题
#
★★★

1. RDD(弹性分布式数据集)的五大特性

Spark 的 RDD(弹性分布式数据集)有哪五大特性,它们分别如何支撑 RDD 的容错与计算?

  • 分区列表(Partitions)
  • 只读分区上的计算函数(Compute)
  • 依赖关系(Dependencies)

RDD(Resilient Distributed Dataset)是 Spark 的核心抽象,五大特性为:①分区列表(Partitions):RDD 被划分为若干个分区,是并行的基本单元,每个分区对应一个任务;②每个分区上的计算函数(Compute):RDD 对每个分区定义了计算函数,把一个分区映射为另一个分区,实现惰性逐分区计算;③依赖关系(Dependencies):记录 RDD 的血缘,包括窄依赖与宽依赖,是容错与阶段划分的依据;④分区器(Partitioner):仅对 Keyed RDD 存在,决定 key 如何分到分区(Hash/Range),影响 shuffle 与同分区内的数据分布;⑤首选位置(Preferred Locations):记录每个分区偏向的节点(如 HDFS 块所在节点),用于调度时实现数据本地性(data locality),减少网络传输。这五特性共同支撑 RDD 的"弹性(可重建)"与"分布式(并行)"。

记忆口诀"分区、计算、依赖、分区器、首选位置"。作答时强调"分区+依赖"是容错与并行核心,"首选位置"体现数据本地性优化。

#
★★★

2. Sort-based Shuffle 原理与调优(shuffle partitions/压缩/溢写)

Spark 的 Sort-based Shuffle 原理是什么,如何调优 shuffle partitions、压缩与溢写?

  • Sort-based Shuffle 的 Map 端排序与聚合
  • shuffle partitions 的影响
  • 压缩与溢写策略

Spark 自 1.2 起默认使用 Sort-based Shuffle(Hash Shuffle 在 2.0 中被移除)。Map 端每个 task 先按 key 在内存中排序(sort),排序后按目标分区(partition)写成溢出文件,所有溢写文件合并成一个分区文件,并生成索引文件记录各分区偏移;Reduce 端按索引fetch自己的分区数据。相比 Hash Shuffle,Sort-based 每个 task 只生成一个数据文件+索引,大幅减少文件数与随机 IO。调优点:spark.sql.shuffle.partitions(默认 200)控制 reduce 端分区数,分区数过少则倾斜、reducer 过载,过多则 shuffle 小文件与调度开销大,应按数据量与执行器核数合理设置;spark.shuffle.compressspark.shuffle.spill.compress 控制压缩,打开可减少网络与磁盘 IO;spark.shuffle.file.bufferspark.reducer.maxSizeInFlight 控制缓冲与传输;spark.shuffle.spill.numElementsForceSpillThreshold 控制溢写时机,避免内存溢出。Shuffle 是 Spark 最重的 IO 环节,通常配合广播 join 减少 shuffle。

核心是"排序+单文件+索引"提升 shuffle 效率。调优围绕"分区数、压缩、缓冲、溢写"四个维度,分区数是最常被问的。

#
★★★

3. Spark 动态资源分配(Dynamic Allocation)的原理与关键参数

Spark 的动态资源分配(Dynamic Allocation)原理是什么,有哪些关键参数,如何启用?

  • 动态分配 Executor 数量
  • 关键参数(shuffle service、min/max executors)
  • 适用场景(流式/多租户)

动态资源分配(Dynamic Allocation)让 Spark 根据任务运行情况动态调整 Executor 数量:任务积压时增加 Executor,空闲超时后回收 Executor,从而在多租户共享集群上避免资源浪费与申请瓶颈。启用需设置 spark.dynamicAllocation.enabled=true,并依赖外部 Shuffle Service(spark.shuffle.service.enabled=true)在 Executor 减少时保留 shuffle 数据,避免已完成的 shuffle 数据丢失。关键参数:spark.dynamicAllocation.minExecutorsmaxExecutors 限定上下界;spark.dynamicAllocation.initialExecutors 初始数量;spark.dynamicAllocation.executorIdleTimeout(空闲回收,默认 60s)与 spark.dynamicAllocation.schedulerBacklogTimeout(积压触发扩容)。它最适合流式作业(数据量随时间波动)与多租户队列,批处理全量作业用固定 executors 更可控。

抓住"积压扩容、空闲回收"的两条触发逻辑,以及"必须配 Shuffle Service"这个前提。回答强调上下界与两种超时参数。

#
★★★

4. Spark Streaming 与 Structured Streaming 的差异,微批 vs 连续处理、端到端 Exactly-Once 如何实现?

Spark Streaming 与 Structured Streaming 有何差异,微批与连续处理如何选,端到端 Exactly-Once 如何实现?

  • DStream(微批)与 DataFrame 流式处理模型的差异
  • 微批 vs 连续处理(Continuous Processing)
  • 端到端 Exactly-Once(端到端幂等 + 事物输出)

Spark Streaming(DStream)基于 RDD 微批,把连续流切成固定间隔的小批次,API 是 RDD 风格,语义统计弱、处理延迟高(秒级);Structured Streaming 基于 DataFrame/DataSet 的流式模型,把流看作"无限表",支持 event-time 窗口、watermark、精确语义,并提供微批(默认)与连续处理(Continuous Processing)两种模式。微批以固定间隔处理,延迟秒级、吞吐高、容错简单;连续处理是真正的低延迟引擎(毫秒级),但支持的操作受限(仅 select/map/filter 等投影类),且不能保证某些精确一次场景。端到端 Exactly-Once 需要"源端 + 计算端 + 输出端"三层配合:源端可重放(如 Kafka 按 offset 精确消费)、计算端用 checkpoint 记录并行处理进度(offset 原子提交)、输出端用幂等写入或事务式输出(如 Kafka sink 的幂等、JDBC 的 foreachBatch + 幂等 upsert),保证任务重试不重复不丢失。

关键区分"模型"(RDD vs DataFrame)与"模式"(微批 vs 连续)。Exactly-Once 是"可重放源 + checkpoint + 幂等/事务输出"的组合,不是单一机制。

#
★★

5. RDD 的血缘(Lineage)与容错机制,为什么窄依赖可以高效重算、宽依赖需要 Checkpoint,与缓存(Cache)的差别如何?

RDD 的血缘(Lineage)如何支撑容错,为什么窄依赖可高效重算而宽依赖需要 Checkpoint,持久化与 Cache 有何差别?

  • Lineage 血缘图与容错重算
  • 窄依赖 vs 宽依赖的恢复差异
  • Checkpoint 与 Cache 的区别

RDD 通过记录血缘(Lineage)支持容错:每个 RDD 记录其父 RDD 与变换函数,当某个分区丢失(如节点故障)时,Spark 沿血缘图重新计算该分区,无需重算整个作业。窄依赖(一个父分区只对应一个子分区)中丢失分区的重算只需重算个别父分区,代价小、可高效恢复;宽依赖(某父分区对应多个子分区,如 shuffle)中丢失分区需重算其所有父分区,代价高,且长血缘链重算成本累积,因此对长依赖链路或关键中间结果应使用 Checkpoint。Checkpoint 把 RDD 物理落盘(通常是 HDFS/对象存储)并截断血缘,重算时直接从磁盘读,避免整条链重算;Cache(持久化)把 RDD 缓存在内存/磁盘,但不截断血缘,节点故障时缓存丢失仍会沿血缘重算。区别:Cache 是"加速重复计算"(血缘仍在),Checkpoint 是"切断链路、彻底容错"(血缘截断)。

核心对比:Cache 保留血缘只加速,Checkpoint 截断血缘保真容错。窄依赖恢复局部、宽依赖恢复全量是选择 checkpoint 的依据。

#
★★

6. 内存管理与 Tungsten 引擎,UnsafeRow 二进制格式与堆外内存管理

Spark 的内存管理与 Tungsten 引擎如何工作,UnsafeRow 二进制格式与堆外内存管理有何优势?

  • Tungsten 引擎的二进制内存表示
  • UnsafeRow 与编码
  • 堆外内存(off-heap)与缓存友好

Tungsten 是 Spark 的底层内存与执行引擎,目标是把内存管理和计算向量化。核心是 UnsafeRow:它把行数据编码为紧凑的二进制格式(固定长度字段按偏移直接寻址,变长字段用 offset+length 数组),避免 Java 对象头的内存开销与 GC 压力,使数据在内存/CPU 缓存中更紧致,并支持 cache line 友好的顺序访问。配合堆外内存(off-heap,配置 spark.memory.offHeap.enabledspark.memory.offHeap.size),数据可绕过 JVM 堆、减少 GC 与行对象拼接,被 shuffle、排序、聚合等算子在本地直接操作。Tungsten 还利用"codegen(代码生成)"把一段逻辑编译成 Java 字节码,减少虚函数调用与对象创建,从而大幅提升吞吐。代价是二进制格式的数据可读性低、需按 schema 解码,且堆外内存需精心管理。

抓住"UnsafeRow 二进制 + off-heap + codegen"三件套。它解决"GC 压力、对象开销、CPU 缓存不友好"。回答强调"空间紧凑、减少 GC、codegen 加速"。

#
★★

7. Catalyst 优化器与 AQE(自适应查询执行),动态合并 shuffle 分区/倾斜处理

Catalyst 优化器与 AQE(自适应查询执行)如何优化 Spark SQL,动态合并 shuffle 分区与倾斜处理是怎么实现的?

  • Catalyst 的规则优化与物理计划优化
  • AQE 的动态合并分区与切换 join 策略
  • AQE 的倾斜处理

Catalyst 是 Spark SQL 的可扩展优化器,经过"分析→逻辑计划优化(谓词下推、列裁剪、常量折叠、剪枝)→物理计划生成(选择 join 策略、数据源接入)→codegen"阶段,把 SQL 优化成高效执行计划。AQE(Adaptive Query Execution,spark.sql.adaptive.enabled=true)是在运行时根据实际 shuffle 数据动态调整物理计划:①动态合并 shuffle 分区(coalescePartitions):当某些 reduce 分区数据量过小时自动合并,减少任务数;②动态切换 join 策略:运行时若某表变小,可把 sort-merge join 切换为 broadcast join;③动态倾斜处理(skewJoin):检测到某个 shuffle 分区数据量远超中位数时,把倾斜分区拆分(salting)成多个 reduce 任务并行处理,避免单任务拖累整体。AQE 让"预先设置的参数"在运行时自适应,显著缓解数据倾斜与分区不均衡。

Catalyst 是"编译期/静态优化",AQE 是"运行期/动态优化",是 Spark 3.x 的核心增强。倾斜处理是 AQE 的杀手锏(动态拆分倾斜分区)。

#
★★

8. 广播 Join 与 Sort-Merge Join 的取舍,broadcast hint 与阈值(spark.sql.autoBroadcastJoinThreshold)

Spark 的广播 Join(Broadcast Join)与 Sort-Merge Join 如何取舍,broadcast hint 与阈值参数如何设置?

  • Broadcast Join 原理与优势
  • Sort-Merge Join 适用场景
  • 阈值与 hint

Broadcast Join 把小表广播到每个 executor 内存,在本地与大规模表进行 map 端 join,避免 shuffle,适合"小表 join 大表"。采用条件:小表大小低于 spark.sql.autoBroadcastJoinThreshold(默认 10MB,-1 表示禁用),或显式 broadcast hintSELECT /*+ BROADCAST(t) */)。Sort-Merge Join 先把两侧表按 join key 排序,再归并连接,适合大表 join 大表,无内存限制但需 shuffle 与排序。取舍:能用 broadcast 就优先(省 shuffle、吞吐高),但小表过大(超过阈值)时广播会撑爆 executor 内存或 GC,此时改用 sort-merge;阈值设置需权衡(小表可装入单 executor 内存)。实际工程中,AQE 会在运行时自动把小表切换为 broadcast join。

关键判断"小表→broadcast、大表→sort-merge"。broadcast 靠"内存换 shuffle",阈值是容量边界。AQE 可自动切换。

#
★★

9. 数据倾斜的定位与治理,salting/skew join hint/两阶段聚合

Spark 数据倾斜如何定位与治理,salting、skew join hint、两阶段聚合分别如何应用?

  • 数据倾斜的定位方法
  • salting 加盐
  • 两阶段聚合与 skew join hint

数据倾斜指某个 key 或分区数据量远超其他,导致个别 task 耗时过长、整体拖慢。定位:查看 stage 中个别 task 处理数据量/耗时远大于其他(Spark UI 的 stage 详情、executor 堆栈),或检查 join/groupBy 的 key 分布。治理手段:①加盐(salting):对倾斜 key 加随机前缀打散,使其分布到多个分区,join 时两侧都加盐(大表加随机盐、小表复制多份),从而把倾斜大任务拆成多个小任务;②skew join hint:Spark 3.4+ 提供 /*+ SKEW(表(k), ((值), (n))) */ hint 或配合 AQE 自动拆分倾斜分区;③两阶段聚合:group by 场景先按"key+随机盐"局部聚合,再去掉盐做全局聚合,避免单 key 集中在一处。本质都是"把倾斜的 key 打散到多个分区并行处理"。

三招对应不同场景:join 倾斜用加盐/skew hint,group by 倾斜用两阶段聚合。核心思想是"打散倾斜 key"。答题先定位再选招。

#
★★

10. Spark DAG 与 Stage/Task 划分,宽依赖(shuffle)与窄依赖(pipeline)的边界

Spark 的 DAG 如何根据宽窄依赖划分 Stage,宽依赖与窄依赖的边界是什么,各自如何执行?

  • 窄依赖与宽依赖的定义
  • Stage 划分规则(宽依赖处切分)
  • 窄依赖的 pipeline 与宽依赖的 shuffle

Spark 依赖分两类:窄依赖(Narrow Dependency)是一个父分区只被一个子分区使用,如 map、filter、union,可在一个 task 内 pipeline 连续执行,不产生 shuffle;宽依赖(Wide Dependency)是一个父分区被多个子分区使用,如 groupByKey、reduceByKey、join 的 shuffle 阶段,需要 shuffle 把数据重新分发。DAG 调度器在"宽依赖处"切分 Stage:从后往前,遇到宽依赖就断开,形成一个新 Stage;窄依赖的算子合并进同一 Stage 内 pipeline 执行。因此一个 Stage 内全是窄依赖(可流水线),Stage 之间是 shuffle 边界。Stage 中的每个分区对应一个 Task,Task 并行度由 Stage 最后一个 RDD 的分区数决定。宽窄依赖划分同时决定了容错重算代价与调度复杂度。

记住"宽依赖一个父喂多个子就要 shuffle,遇到宽依赖切 Stage"。窄依赖可 pipeline 是提升性能的关键。这是 DAG 调度的核心。

#
★★

11. Spark Job/Stage/Task 的划分,DAG 调度器如何根据宽窄依赖切分 Stage,Stage 内并行度由谁决定?

Spark 的 Job/Stage/Task 如何划分,DAG 调度器如何根据宽窄依赖切分 Stage,Stage 内并行度由谁决定?

  • Job/Stage/Task 三级划分
  • DAG 调度器按宽依赖切分 Stage
  • Stage 内并行度由分区数决定

Spark 将一个应用(Application)划分为多个 Job(每个 Action 触发一个 Job),每个 Job 对应一个 DAG,DAG 被 DAG 调度器按宽依赖切分为多个 Stage,每个 Stage 的每个分区对应一个 Task 由 TaskScheduler 调度。切分规则:从最后一个 RDD 沿依赖往前遍历,遇到宽依赖(shuffle)就断开,形成新的 Stage 边界;因此 Stage 内只含窄依赖,可 pipeline 执行。Stage 内并行度(Task 数量)由该 Stage 最后一个 RDD 的分区数决定(即 shuffle 后的分区数,可通过 spark.sql.shuffle.partitions 控制),而非 Executor 数量;Task 被分配到空余 Executor 上调度执行。一个 Job 的 DAG 由一系列 Stage 组成,Stage 之间靠 shuffle 串联。

记忆"Job 由 Action 触发、Stage 由宽依赖切分、Task 由分区数决定"。Stage 内并行度=最后一个 RDD 分区数,这是容易被忽略的点。

#
★★

12. Spark 的数据倾斜治理全览,定位(stage 耗时/堆栈)、加盐、两阶段聚合、skew join hint、动态分区重分布如何?

Spark 数据倾斜治理的完整方案是什么,如何定位,加盐、两阶段聚合、skew join hint、动态分区重分布分别如何用?

  • 定位(stage 耗时、堆栈、趋势)
  • 加盐与两阶段聚合
  • skew join hint 与动态分区重分布

数据倾斜治理全流程:先定位再治理。定位:Spark UI 中某 stage 的多个 task 耗时差异巨大、个别 task 长时间卡住,或堆栈显示某 task 处理的数据量远超中位数,可进一步对 join key 做 count 分布统计。治理手段梯度:①加盐(salting):对倾斜 key 加随机前缀打散到多分区,join 时两侧配合(大表随机盐、小表复制多份);②两阶段聚合:先按 key 加盐局部聚合再全局聚合,适用于 group by/聚合倾斜;③skew join hint:/*+ SKEW(...) */ 显式指定倾斜 key 及其分区数,Spark 拆分倾斜分区单独处理;④动态分区重分布(AQE 的 spark.sql.adaptive.coalescePartitions 与 skew 优化):运行时自动把倾斜分区拆成多个子分区并行;⑤增加 shuffle 分区数或调整并行度。原则是"先小改动(参数/分区)再针对性打散,优先使用 AQE 自动处理"。

治理不是单一手段,而是"定位→选择"的流程。AQE 自动倾斜处理是首选,手动加盐/两阶段聚合用于兜底。回答体现"从自动到手动"的梯度。

#
★★

13. Spark 的调度,DAG→Stage→Task 划分,宽窄依赖与 stage 间 shuffle 如何?

Spark 的调度机制如何从 DAG 划分到 Stage 再到 Task,宽窄依赖与 stage 间 shuffle 的关系是什么?

  • DAG 调度器与 Task 调度器
  • 宽窄依赖与 shuffle 边界
  • 调度执行流程

Spark 调度分两层:DAG 调度器(DAGScheduler)负责高层调度,把 Job 的 DAG 按宽依赖切分成 Stage,并处理 Stage 的依赖关系(树状拓扑),为每个 Stage 生成 Task 集提交给 Task 调度器;Task 调度器(TaskScheduler)负责把 Task 分发到 Executor 执行、处理失败重试与推测执行。宽依赖(shuffle)是 Stage 之间的天然边界:一个宽依赖产生一次 shuffle,把上游 Stage 的输出按 key 重新分区给下游 Stage;窄依赖算子合并进同一 Stage 内 pipeline 执行,无 shuffle。因此"Stage 间有 shuffle、Stage 内无 shuffle"是 Spark 调度与执行的核心特征。调度时还会结合数据本地性(首选位置)把 Task 优先调度到数据所在节点。

分清 DAGScheduler(切 Stage)与 TaskScheduler(跑 Task)两层职责。宽依赖=shuffle 边界=Stage 边界,是记忆锚点。

#
★★

14. RDD 持久化级别(StorageLevel)的选择与 Kryo 序列化配置

Spark 的 RDD 持久化级别(StorageLevel)如何选择,Kryo 序列化如何配置,各有什么权衡?

  • StorageLevel 的内存/磁盘/序列化组合
  • 各级别的选择依据
  • Kryo 序列化配置与优势

StorageLevel 定义持久化方式,由"存储介质(内存/磁盘)+ 是否序列化 + 副本数"三维组合而成,常用级别:MEMORY_ONLY(只存内存、对象形式,最快但内存不足会丢分区)、MEMORY_AND_DISK(内存放不下溢写磁盘,容错好)、MEMORY_ONLY_SER / MEMORY_AND_DISK_SER(序列化后存储,省内存但增加 CPU 与反序列化开销)、以及 _2 后缀(副本数 2,提升容错)。选择原则:数据小且会多次复用→MEMORY_ONLY;数据大内存不足→MEMORY_AND_DISK_SER;内存紧张→用序列化;对可靠性要求高→加副本。Kryo 序列化比 Java 默认序列化快得多、体积小,配置 spark.serializer=org.apache.spark.serializer.KryoSerializerspark.kryo.registrationRequired 注册类可进一步优化,适合 shuffle 与持久化的序列化数据。Trade-off:Kryo 需注册类、不友好于未注册自定义类,但换性能。

记住 StorageLevel 三维:内存/磁盘、序列化/非序列化、副本数。Kryo 是"更快更省但需注册"的序列化。选择是"内存空间 vs CPU 开销 vs 容错"权衡。

#
★★

15. Spark SQL 的谓词下推与列裁剪,DataSource V2 的过滤器下推

Spark SQL 的谓词下推与列裁剪如何优化 IO,DataSource V2 如何实现过滤器下推?

  • 谓词下推(Predicate Pushdown)
  • 列裁剪(Column Pruning)
  • DataSource V2 的过滤下推接口

谓词下推(Predicate Pushdown)把 WHERE 过滤条件尽量下推到数据源或扫描阶段,只读取满足条件的行,减少网络与 IO;列裁剪(Column Pruning)只读取查询需要的列,避免读取整行数据,两者配合能显著减少扫描数据量。这些优化由 Catalyst 在逻辑计划阶段应用,对文件源(Parquet)可下推谓词与列,对 Hive 表下推分区裁剪。DataSource V2(DataSourceV2)提供更完善的过滤器下推:通过 SupportsPushDownFilters 接口把表达式过滤(如 IsNullEqualToGreaterThan 等)下推给数据源,数据源在读取时用这些过滤器裁剪分区/行/文件,减少扫描;同时通过 SupportsPushDownCatalystFilters 支持更复杂谓词。V2 相比 V1 有统一的 data source 读写接口、支持谓词/列/分区裁剪下推、流式等能力,是 Spark 读取外部数据源(如 Iceberg、JDBC、文件)的推荐途径。

记忆"谓词下推=少读行、列裁剪=少读列",二者是 Catalyst 的核心 IO 优化。DataSource V2 把过滤下推给数据源进一步提升裁剪。

#

16. Spark 与 Flink 的选型,批流一体、状态管理、延迟与吞吐在不同场景下的取舍如何?

Spark 与 Flink 在批流一体、状态管理、延迟与吞吐方面如何选型取舍?

  • 批流一体能力对比
  • 状态管理与窗口语义
  • 延迟与吞吐取舍

Spark 与 Flink 是现代大数据两大引擎,选型看场景:Spark 以"批处理"见长,微批流式(Structured Streaming)吞吐高、生态成熟(SQL/ML/图),但延迟秒级、状态管理相对弱(基于 checkpoint 的状态);Flink 以"流处理"见长,真流式(逐事件)延迟毫秒级,原生支持高效状态管理(可配置 TTL、精确一次)、event-time 与高级窗口,且 Flink 1.12+ 实现了批流一体(同一套 API 跑批与流)。取舍:对延迟不敏感、重批处理、依赖 Spark 生态的离线场景选 Spark;对实时性、复杂状态、精确一次、低延迟要求高的场景选 Flink;批流一体且需要统一 API 的场景越来越多选 Flink,但 Spark 的 SQL 与批处理生态仍具优势。两者也常组合:Spark 跑离线批、Flink 跑实时流,共享湖仓(Lakehouse)。

用"延迟 vs 吞吐"与"批 vs 流"两根轴定位:Spark 偏批、吞吐高、延迟高;Flink 偏流、延迟低、状态强。回答要落到业务场景。

#

17. Spark 内存模型,执行内存与存储内存的 Unified Memory 管理如何?

Spark 的内存模型如何管理执行内存与存储内存,Unified Memory 的统一管理机制是什么?

  • 堆内存与堆外内存分区
  • Unified Memory 的弹性借用
  • 内存参数(spark.memory.fraction 等)

Spark 使用 Unified Memory(统一内存)模型管理 JVM 堆内存,分为三块:Reserved(保留,约 300MB)、User Memory(用户内存,供用户代码/数据结构使用)、Unified Memory(执行+存储共享)。Unified Memory 内部为"执行内存(execution,用于 shuffle、join、聚合、sort)"与"存储内存(storage,用于缓存 RDD/广播变量)"的动态分区:两者可互相借用(storage 可借用 execution 的闲置内存,execution 也可抢占 storage 已占用的内存,但 execution 在需要时强制驱逐 storage 的已缓存数据),从而提升内存利用率。参数:spark.memory.fraction(默认 0.6,Unified 占堆比例)、spark.memory.storageFraction(默认 0.5,storage 初始占比)。堆外内存(off-heap)由 spark.memory.offHeap.enabledspark.memory.offHeap.size 控制。执行内存不足会导致频繁溢写磁盘,storage 过分占用会挤压执行内存,需平衡。

核心是"执行与存储柔性共享、动态互借"。两个 fraction 参数是常考点。注意 execution 有抢占权,storage 可被驱逐。

#

18. Spark 的 checkpoint 与 cache 的区别,血缘关系与容错恢复的差异如何?

Spark 的 checkpoint 与 cache 有什么区别,对血缘关系与容错恢复有何不同影响?

  • Checkpoint 落盘并截断血缘
  • Cache 保留血缘只加速
  • 容错恢复路径差异

Cache(持久化)与 Checkpoint 都用于保存 RDD 中间结果,但本质不同:Cache 把 RDD 缓存在内存(或磁盘)以加速重复计算,但保留血缘,一旦缓存数据因节点故障丢失,Spark 会沿血缘图重新计算该 RDD;Checkpoint 把 RDD 物理写入可靠存储(HDFS/对象存储)并截断血缘(斩断父 RDD 依赖),数据由外部存储保证,故障时直接从存储恢复而不重算祖先链。使用场景:计算链长、重算代价高或需彻底容错的中间结果用 Checkpoint;同一 RDD 被多次 Action 复用且重算代价低时用 Cache。Checkpoint 建议在 checkpoint 前先 cache(避免重复计算写盘)。代价:Checkpoint 落盘开销大,Cache 内存开销大。

一句话区分"Checkpoint 断血缘、Cache 留血缘"。Cache 是性能优化,Checkpoint 是容错机制。恢复路径不同:Cache 回放血缘,Checkpoint 读存储。

#

19. 广播变量(Broadcast)的实现与内存边界

Spark 的广播变量(Broadcast)如何实现,其内存边界是什么,如何合理使用?

  • Broadcast 的机制(数据只发一次、各 task 共享)
  • 内存边界与阈值
  • 使用场景与注意事项

广播变量(Broadcast)把一份只读数据发送到每个 Executor 一次,供该 Executor 上所有 Task 共享,避免每个 task 都序列化拷贝一份数据(否则大表 join 时数据被重复传输、撑爆内存)。实现上,Driver 把数据序列化后通过 Torrent 协议(分块点对点)分发到各 Executor,各 Executor 缓存该广播变量,Task 用 broadcast.value 读取。内存边界:广播数据常驻 Executor 内存(属于 storage 内存),若广播数据过大(超过阈值或 Executor 内存),会撑爆内存、GC 或 OOM;spark.sql.autoBroadcastJoinThreshold 默认 10MB 控制 SQL 自动广播的阈值,手动 broadcast() 则需自担内存风险。使用建议:只广播小到中等大小的只读数据(如维表、配置),避免广播大表;用 Kryo 压缩减小体积;用完可 unpersist 释放。

广播的本质是"一次分发、共享复用、以内存换传输"。内存边界是"广播数据常驻 Executor 内存",数据过大是风险点。回答要强调"适合小数据"。