数据岗位面试题更新 2026-08-05

请说明在 PySpark 中调用 checkpoint() 对中间结果执行检查点操作的具体用法与适用场景。

数据技术原理方案权衡PySpark

考察说明

考查对 PySpark checkpoint 机制及其在容错与优化中作用的理解。

回答思路

  1. 【回答框架 1】checkpoint() 是 RDD 与 DataFrame 提供的容错机制,通过将中间计算结果落盘(本地或 HDFS)以截断 lineage,避免依赖链过长导致重算代价高。
  2. 【回答框架 2】用法上,调用前需设置 checkpoint 目录,如 sc.setCheckpointDir('hdfs://...') 或 spark.sparkContext.setCheckpointDir,然后对中间结果调用 checkpoint();DataFrame 还需注意 checkpoint 是 action 操作,会触发一次实际计算。
  3. 【回答框架 3】checkpoint 与 cache 不同:cache 将数据存储在内存或磁盘,保留 lineage 以便故障时重建;checkpoint 直接切断血缘,将数据作为新起点,适合迭代或复杂 DAG 中减少恢复与重算开销。
  4. 【回答框架 4】注意 checkpoint 是惰性执行,需在后续 action 时才真正落盘;多次 checkpoint 会引入额外 I/O 开销,应只在关键中间结果或长依赖链时使用。
  5. 【关键点 1】checkpoint() 用于截断 RDD/DataFrame 的血缘链,提高容错与恢复效率。
  6. 【关键点 2】使用前必须设置检查点目录,否则运行报错。
  7. 【关键点 3】checkpoint 会触发实际计算,与 cache 的惰性语义不同。
  8. 【关键点 4】适合迭代算法或 DAG 过长的场景,避免故障时全链重算。
  9. 【易错点 1】混淆 checkpoint 与 cache,误以为 checkpoint 保留血缘并可快速复用。
  10. 【易错点 2】忽略 checkpoint 的 action 触发特性,导致性能分析失误。
  11. 【易错点 3】频繁 checkpoint 增加 I/O 与存储压力,需权衡使用频率。