请从任务调度的角度解释 Spark Streaming 如何利用 DAG(有向无环图)来执行作业,包括阶段划分与任务提交的流程。
考察说明
考查对 Spark Streaming 任务执行机制的理解,特别是 DAG 在作业调度中的角色。
回答思路
- 【回答框架 1】Spark Streaming 基于微批处理模型,将连续数据流切分为一系列 RDD 批次,每个批次对应一个作业。这些 RDD 之间通过依赖关系构成 DAG,描述了计算步骤的拓扑结构。
- 【回答框架 2】作业提交时,DAGScheduler 会将 DAG 中的 RDD 根据宽依赖(如 shuffle 依赖)划分为多个 Stage,每个 Stage 包含一组可并行执行的任务。窄依赖的转换会在同一 Stage 内流水线执行。
- 【回答框架 3】Stage 按照依赖顺序依次提交,每个 Stage 被分解为多个 Task,Task 被调度到 Executor 上执行。默认使用 FIFO 或 Fair 调度器,任务执行结果会反馈给 DAGScheduler。
- 【回答框架 4】DAG 的执行保障了作业的有效顺序与容错性。若某个分区数据丢失,可根据血缘关系重新计算,只重算丢失分区的部分。
- 【关键点 1】Spark Streaming 将数据流按批次转换为 RDD,构建作业 DAG。
- 【关键点 2】DAGScheduler 依据宽依赖划分 Stage,窄依赖在同一 Stage 内流水线化。
- 【关键点 3】Stage 被拆分为 Task 并分发到 Executor 执行。
- 【关键点 4】DAG 血缘关系支持容错,无需完全重算。
- 【易错点 1】不要混淆 Spark Streaming 的微批处理与纯流处理,DAG 是针对每个批次的作业。
- 【易错点 2】切勿误以为同一 Stage 内的任务完全独立,它们可能共享分区数据。
- 【易错点 3】不要忽略调度策略对任务执行顺序的影响。