适用范围与学习目标
本文以 Apache Spark 3.5.7 为基线,适用于理解批处理和 Spark SQL 应用的执行过程。YARN、Kubernetes 与 Standalone 的资源管理实现不同,本文介绍共同执行概念,并用 YARN 作为运维对照;Spark 4.x 和商业发行版的依赖、计划及配置应重新核对。示例未在实际集群运行,验证状态为 PENDING。
目标是解释 Driver、Executor、Job、Stage、Task 和 Shuffle 之间的关系,并把应用排队、计算变慢与结果写入失败分开。Spark 是计算引擎,数据可以来自 HDFS、对象存储或数据库;看到 Spark 进程并不能推断数据一定在内存中,也不能推断集群采用哪种存储。
Driver 与 Executor 各自做什么
Driver 运行应用主程序并协调任务,Executor 在集群节点中执行任务、存储缓存和 Shuffle 数据。资源管理器负责为应用提供执行资源,不负责替代 Driver 的 SQL 计划与任务调度。不同应用通常有各自的 Executor,不能把某应用缓存的 DataFrame 视为所有应用共享的数据缓存。
部署模式决定 Driver 位于提交客户端还是集群内部。client 模式下提交机器的网络和进程稳定性是应用依赖;cluster 模式下应去相应的集群日志系统找 Driver 输出。Executor 日志不能替代 Driver 异常记录。职责与部署模式依据 Spark 3.5.7 集群模式概览。
从惰性计算理解 Job、Stage 与 Task
DataFrame 或 RDD 的转换通常先形成计算描述,真正需要结果的 action 才触发执行。一个应用可以产生多个 Job,每个 Job 由若干 Stage 组成,Stage 内的 Task 面向数据分区并行执行。应区分应用、SQL execution、Job 和 Stage 的标识,不能把同一个 SQL 中的多次 Job 误认为重复提交应用。
某些转换能在已有分区内继续计算,另一些操作需要重新分布数据。Shuffle 往往成为阶段边界并带来磁盘、网络和重试成本。多给 CPU 不一定能加速分区数量很少或数据严重倾斜的 Stage。惰性计算及 Shuffle 的概念见 RDD Programming Guide。
用一千行样本观察计划与结果
下面是使用现有 SparkSession 的 PySpark 3.5.7 教学示例,只构造一千个整数,不读取业务表、不写入外部系统。执行 show 会提交实际计算,所以在学习环境中使用:
from pyspark.sql import functions as F
sample = spark.range(0, 1000, 1, 4)
summary = sample.groupBy((F.col("id") % 10).alias("bucket")).count()
summary.explain(mode="formatted")
summary.orderBy("bucket").show(10, truncate=False)若示例成功,数学上应有十个 bucket,每组 count 为一百;这是预期,不是本文取得的实测输出。物理计划通常包含聚合与 Exchange 等节点,具体名字、层次和最终分区数会受版本和 AQE 配置影响,不应以计划文本逐字相同为验收标准。orderBy 也可能引入额外排序过程,因此不能据此例给复杂 SQL 估算固定 Stage 数。
SQL 计划与运行时计划要分开看
SQL 会经历解析、分析、优化和物理计划选择;字段能解析不代表数据一定可读,物理计划能生成也不代表运行一定成功。用 EXPLAIN FORMATTED 或 DataFrame explain 理解扫描、过滤、关联和聚合的位置,再对照运行中的 SQL 页面和实际指标。
Spark 3.5.7 的 AQE 可以利用运行时统计调整计划,因此提交前的解释结果不一定等于最终执行方案。若要证明倾斜优化或广播转换生效,查看已完成 SQL 的最终计划与指标。格式和模式定义见 EXPLAIN 文档,参数边界见 SQL 调优指南。
用一条时间线定位应用慢在哪里
- 提交到 Driver 启动:优先证据:调度器、AM 或 Pod 状态;典型候选原因:队列、权限、资源申请。
- Driver 启动到首批 Task:优先证据:Driver 日志、文件发现与计划;典型候选原因:元数据访问、依赖加载、计划开销。
- Task 执行期间:优先证据:Stage 分布、Executor 指标;典型候选原因:扫描、Shuffle、倾斜、GC、外部服务。
- Task 完成到结果可见:优先证据:输出提交日志、目标系统状态;典型候选原因:提交协议、重试、目录或事务发布。
采样时保存 application ID、attempt、SQL execution ID 和时区。计算结束后的业务延迟可能来自下游目录刷新或索引,不能只按 Spark UI 的完成时间结案。YARN 中尚未取得资源的问题,接着使用 YARN 排队诊断。
缓存、重算与失败重试的边界
缓存适合复用成本高、重复访问且资源允许的数据。首次访问仍需要填充缓存;资源不足或 Executor 丢失时,分区可能需要重新计算。缓存不等于备份,应用结束后的持久保存需要明确外部存储与提交结果。盲目缓存所有中间表会与执行内存竞争,应测量重复利用次数及缓存命中情况。
Task 重试意味着同一段用户代码可能再次执行,写数据库、调用接口等外部副作用需要单独保证幂等或事务语义。普通应用重启也不自动等于 Structured Streaming 从原 checkpoint 恢复;两者的输入位点与状态管理不同。流处理可继续阅读 Spark 流任务 checkpoint 与恢复。
常见误区与观测准备
把所有错误都归因于 Executor 内存不足,会忽略 Driver 上的大规模 collect、文件列表和计划开销。collect 把结果带回 Driver,结果规模未知时不能拿它做生产数据预览。另一个误区是将更多 Executor 视为通用加速器;输入分区、倾斜、外部吞吐和资源队列都可能限制实际收益。
为事后分析配置受控的 event log 保存与 History Server,并确认日志路径权限、保留周期及敏感字段处理。History Server 依赖事件日志,不能保证每个未启用日志或日志缺失的应用都能完整回放。监控能力依据 Spark 3.5.7 Monitoring。
验收、停止与回退边界
学习验收应能从样本计划指出一次数据重分布,解释预期分组结果,并把一个真实应用的等待、计算和提交阶段分开。运维验收还要保存数据窗口、有效配置、日志来源与阶段耗时,不把示例成功当作生产性能证明。后续可用 Spark SQL 性能诊断 形成一次可比较优化实验。
本文只读计划没有业务数据回退步骤;教学 show 只计算有限数据。若在真实应用调优中产生新输出,应写到独立候选目标再验收。取消或重跑应用之前,先检查已经提交的外部写入与幂等边界;进程停止不会撤销外部系统中已成功的写入。