DAG 执行框架优于 MapReduce 的地方在哪里?
Created|Updated|系统架构
|Word Count:210|Reading Time:1mins|Post Views:
有个同学问我什么是 DAG 框架。我感觉隐隐约约听过,但又讲不清楚它的概念。
上网搜了一下,我们常见的新大数据执行框架如 Spark、Storm,还有一个我没听过的 Tez,都算 DAG 任务执行框架。他们的主要优点是,可以用 DAG 事先通晓整个任务的全部步骤,然后进行转换优化。如 Tez 就可以把多个任务转换为一个大任务,而 Spark 则可以把相关联的 Map 直接串联起来, 免得多次写回 hdfs(看来 hdfs 也很慢)。传统的 MapReduce 框架为什么不能理解这种优化空间的存在,在任务运行的时候好像一个盲人一样,是个很有意思的话题。
Quora 上的一个相关的问答。
Author: magicliang
Copyright Notice: All articles on this blog are licensed under CC BY-NC-SA 4.0 unless otherwise stated.
Related Articles
2026-07-13
宽依赖与窄依赖:Stage 是怎样划分出来的
上一篇拆解了 RDD 的五大属性,其中 dependencies 决定了 RDD 之间的依赖类型。这一篇进入依赖类型的核心区分——宽依赖和窄依赖——以及它如何直接决定 Stage 的划分。 宽依赖和窄依赖容易被理解成"一对一"和"多对多"的分区关系。更准确的说法是:窄依赖意味着子 RDD 的每个分区只依赖父 RDD 的固定少数分区,可以在单个 Task 内完成计算;宽依赖意味着子 RDD 的分区需要读取父 RDD 所有分区的数据,必须等待一次全量数据交换(Shuffle)。 本文只抓一个问题:DAGScheduler 按什么规则把 RDD 依赖图切分成 Stage。 下图展示了 Stage 划分的核心规则——窄依赖 pipeline 执行,宽依赖切出新 Stage: 窄依赖的三种形态 窄依赖(NarrowDependency)的定义是:父 RDD 的每个分区最多被子 RDD 的一个分区使用。在 Spark 源码中,NarrowDependency 有两个具体子类: OneToOneDependency:父子分区一一对应。map、filte...
2026-07-13
动态资源分配与作业调度:多租户集群的资源博弈
上一篇解决了 Spark 在 YARN 和 K8s 上的部署模式。这一篇进入资源调度——在多个作业共享集群时,Spark 如何动态分配和回收资源。 默认情况下,一个 Spark 应用在启动时申请固定数量的 Executor,作业结束前一直占用。这在单用户开发环境下可以接受,在多租户生产集群中会造成严重的资源浪费——一个空闲的 Spark 应用占着 100 个 Executor 不释放,其他作业排队等待。 本文只抓一个问题:动态资源分配如何按需申请和释放 Executor,以及 FAIR 调度器如何在多个作业之间公平分配资源。 动态资源分配(DRA) Dynamic Resource Allocation(DRA)允许 Spark 应用根据工作负载动态增减 Executor 数量:有 pending Task 时申请新 Executor,Executor 空闲超时后释放。 12345678910111213141516作业生命周期中的 Executor 数量变化: Executor 数量 8 │ ┌────┐ 6 │ ┌────┤ ├──...
2026-07-13
DataFrame 与 Dataset:类型安全与性能的折中
上一篇解决了 Catalyst 优化器的四阶段流水线。这一篇进入用户 API 层——RDD、DataFrame、Dataset 三代 API 的演进逻辑。 三者的区别容易被简化为"RDD 是低级 API,DataFrame 是高级 API"。更准确的说法是:三代 API 在类型安全和执行优化之间做了不同的取舍。RDD 完全类型安全但无法被 Catalyst 优化;DataFrame 完全被 Catalyst 优化但放弃了编译期类型检查;Dataset 试图兼顾两者,但付出了 Encoder 序列化/反序列化的代价。 本文只抓一个问题:三代 API 在内部表示、优化路径和序列化机制上的具体差异。 三代 API 的内部表示 1234567891011121314RDD[Person] └─ 内部存储: JVM 堆上的 Java/Scala 对象 └─ 优化路径: 无(用户代码是黑盒,Spark 无法查看函数内部) └─ 类型信息: 编译期完整保留(泛型参数 T = Person)DataFrame (= Dataset[Row]) └─ 内部存储: Uns...
2026-07-13
DAG Scheduler 与 Task Scheduler:从逻辑计划到物理执行
上一篇解决了 Stage 的划分规则——遇到 ShuffleDependency 就切一刀。这一篇进入调度器内部。 Stage 划分出来之后,谁来决定 Stage 的提交顺序?谁来把 Stage 拆成 Task 发给 Executor?Spark 用两层调度器分工完成这件事:DAGScheduler 负责 Stage 级别的依赖分析和提交顺序,TaskScheduler 负责 Task 级别的资源分配和执行调度。 本文只抓一个问题:一个 action 触发之后,从 Job 到 Stage 到 Task 再到 Executor,调度链路上每一步发生了什么。 调度全景 12345678910111213141516171819202122232425用户代码: rdd.count() │ ▼ SparkContext.runJob() │ ▼ DAGScheduler.submitJob() │ 构建 Stage DAG,按依赖顺序提交 ▼ DAGScheduler.submitStage() │ 检查父 S...
2026-07-13
Structured Streaming:把流当成无界表
上一篇解决了数据源 API 和谓词下推。这一篇进入流处理——Structured Streaming。 Spark 的流处理经历了两代:DStream(基于 RDD 的离散化流)和 Structured Streaming(基于 DataFrame/Dataset 的结构化流)。DStream 已经停止演进,Structured Streaming 是当前唯一推荐的流处理 API。 Structured Streaming 的核心模型可以用一句话概括:把流数据看成一张不断追加新行的无界表,复用 Spark SQL 的 Catalyst 优化器和 Tungsten 执行引擎来处理。 本文只抓一个问题:无界表模型如何工作,以及 MicroBatch 执行引擎如何把连续的流切成离散的批次。 无界表模型 12345678传统流处理的心智模型: 消息 → 处理函数 → 输出 (一条一条处理,每条消息触发一次计算)Structured Streaming 的心智模型: 数据源不断追加行到一张"输入表" 查询在整张表上持续运行 结果写入一张"输出表&qu...
2026-07-13
导读:为什么 Spark 的核心是一张 DAG
Spark 常被介绍为"比 MapReduce 快 100 倍的计算引擎"。这个说法把重点放错了地方。速度差异只是结果,不是原因。真正的区别在于执行模型:MapReduce 是一条固定的 Map-Shuffle-Reduce 三段流水线,而 Spark 是一个可以表达任意有向无环图(DAG)的执行引擎。 本篇只抓一个问题:Spark 的 DAG 执行模型到底和 MapReduce 的线性管道有什么本质区别,以及这个区别如何贯穿整个系列。 MapReduce 的线性管道 MapReduce 的执行模型可以用一句话概括:读数据 → Map → 写磁盘 → Shuffle → 写磁盘 → Reduce → 写数据。 1HDFS → Map Task → Local Disk → Shuffle → Local Disk → Reduce Task → HDFS 这条管道有两个结构性约束: 第一,每个 MapReduce 作业只有一轮 Map 和一轮 Reduce。如果业务逻辑需要多轮聚合(比如先按用户分组,再按地区汇总),就必须把逻辑拆成多个 MapReduce 作...
Announcement
人生只是,守株待兔

