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
Spark on YARN / K8s:资源管理与部署模式
上一篇解决了流处理的时间语义。这一篇进入部署架构——Spark 如何在不同的集群管理器上运行。 Spark 的一个设计决策是把计算引擎和资源管理解耦。SparkContext 通过 SchedulerBackend 接口与集群管理器交互,不绑定任何特定的资源管理系统。这让同一份 Spark 代码可以运行在 Standalone、YARN、Kubernetes 或 Mesos 上,只需要切换启动参数。 本文只抓一个问题:Client Mode 和 Cluster Mode 的区别,以及 YARN 和 Kubernetes 两种主流部署模式各自的架构特点。 Driver 和 Executor 的角色 123456789101112131415┌──────────────────────────────────────────────────┐│ Driver (SparkContext 所在进程) ││ ├── DAGScheduler: RDD → Stage ││ ├── TaskScheduler...
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
RDD:弹性分布式数据集的五大属性
上一篇确立了 Spark 的核心是一张 DAG,而 DAG 的节点就是 RDD。这一篇进入 RDD 本身。 RDD 容易被理解成"分布式的数组"或者"分布在多台机器上的数据集合"。更准确的说法是:RDD 是一份计算配方,记录了数据从哪来、经过什么变换、丢失一个分区后怎么重算。 本文只抓一个问题:RDD 的五大属性分别控制了什么,以及这五个属性如何支撑 Spark 的调度和容错。 五大属性总览 Spark 源码中 RDD 的抽象类定义了五个方法,每个方法对应一个核心属性: 123456RDD[T]├── getPartitions: Array[Partition] ← 数据怎么切分├── getDependencies: Seq[Dependency[_]] ← 上游是谁├── compute(split, context): Iterator[T] ← 一个分区怎么算├── partitioner: Option[Partitioner] ← 按什么规则分区└── getPreferredLocations(s...
2026-07-26
深入 Hadoop 00 - 导读:节点总会失败
Hadoop 常被介绍为"HDFS + YARN + MapReduce 三驾马车"。这个分类把重点放错了位置。把存储、调度、计算切成三件独立的事情,读者很难回答下面这个问题:为什么这三个子系统会一起出现在同一个项目里,而不是像 Spark、Flink、Kafka 那样作为独立组件存在? 更准确的说法是:Hadoop 是一组关于"如何构建运行在大量廉价节点上的分布式系统"的前提假设在三个不同层面的具体化。这组前提把整个项目串成一条主线——节点总会失败、数据总是巨大、廉价比专用更重要。HDFS、YARN、MapReduce 各自处理这组前提下的一个子问题:HDFS 把大文件切成块并多副本放置,YARN 把集群资源切成容器并按策略分配,MapReduce 把大计算切成任务并按失败重试执行。 本系列只抓一个问题:Hadoop 这套以"节点总会失败"为前提假设的工程化方案,在 HDFS、YARN、MapReduce 三个层面是如何具体实现的,每个实现里有哪些可以迁移到其他分布式系统的设计模式。 GFS 论文留下的设计遗产 要理解...
2026-07-13
数据源 API 与谓词下推:让存储层帮忙过滤
上一篇解决了 AQE 的运行时优化。这一篇进入数据源层——Spark 如何把计算推到存储层执行。 查询优化不只发生在执行引擎内部。如果存储层能在读取数据时就过滤掉不需要的行和列,引擎需要处理的数据量会大幅减少。Spark 通过 DataSource API 和 Catalyst 的优化规则把这个能力标准化了:谓词下推(Predicate Pushdown)让存储层只返回满足条件的行,列裁剪(Column Pruning)让存储层只返回查询需要的列。 本文只抓一个问题:谓词下推和列裁剪如何从 Catalyst 的逻辑计划传递到数据源实现。 DataSource API 的两代演进 Spark 的数据源 API 经历了两代: 12345678910DataSource V1 (Spark 1.3+): 接口: InputFormat / OutputFormat + createRelation 谓词下推: 通过 PrunedFilteredScan trait 可选实现 问题: API 不够灵活,不支持流式读写,难以做细粒度优化DataSource V2 (Spark 2.3...
2026-07-13
Adaptive Query Execution:运行时改写执行计划
上一篇解决了 DataFrame 和 Dataset 的 API 演进。这一篇进入运行时优化——Adaptive Query Execution(AQE)。 传统查询优化器的一个根本问题是:优化发生在执行之前,依赖的统计信息可能不准确或完全缺失。AQE 把优化推迟到执行过程中——在每个 Shuffle 边界,收集实际的数据统计信息,然后用这些精确的统计重新优化后续 Stage 的执行计划。 本文只抓一个问题:AQE 的三个核心优化分别解决什么问题,以及运行时优化的反馈循环如何工作。 运行时反馈循环 AQE 的工作方式可以用一句话概括:执行一个 Stage → 收集 Shuffle 输出的统计信息 → 用统计信息重新优化下一个 Stage 的计划 → 执行下一个 Stage。 123456789101112传统执行: 编译期确定完整计划 → Stage 0 → Stage 1 → Stage 2 (统计信息可能不准,计划一旦确定不可修改)AQE 执行: 编译期确定初始计划 → 执行 Stage 0 → 收集 Stage 0 的 Shuffle 输出统计 →...
Announcement
人生只是,守株待兔
