datawarehouse相关
Created|Updated|工程实践
|Word Count:0|Reading Time:1mins|Post Views:
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 与 Flink:两种计算引擎的设计选择
上一篇解决了 Spark 的性能诊断方法。这一篇进入设计对比——Spark 和 Flink 两种计算引擎的架构选择。 这不是一个"谁更好"的问题。Spark 和 Flink 从不同的起点出发,做了不同的设计取舍,各自在擅长的场景中有结构性优势。Spark 从批处理出发,用 micro-batch 扩展到流处理;Flink 从流处理出发,把批看作有界流。两条路线在功能上逐步趋同,但底层架构的差异导致了性能特性和适用场景的持久分化。 本文只抓一个问题:两个引擎在执行模型、状态管理、容错机制和延迟特性上的具体差异及其设计原因。 执行模型对比 Spark 的执行模型是 Stage-based DAG:把计算图按 Shuffle 边界切成 Stage,每个 Stage 包含一组可以 pipeline 执行的窄依赖变换。Stage 之间是全局同步点——上游 Stage 的所有 Task 必须全部完成后,下游 Stage 才能启动。 1234567891011Spark 执行模型: Stage 0: [Task 0] [Task 1] [Task 2] 全部完成 ...
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-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
容错机制:Lineage、Checkpoint 与推测执行
上一篇解决了 BlockManager 的存储管理。这一篇进入容错机制——Spark 如何在节点故障时恢复计算。 分布式系统的容错通常有两条路线:数据副本(每份数据存多份,一份丢了从副本恢复)和计算重放(不存副本,丢了从头重算)。Spark 选择了第二条路线——通过 RDD 的 Lineage(血统)记录计算路径,丢失的分区沿着 Lineage 重算。 本文只抓一个问题:Lineage 容错的优势和代价,以及 Checkpoint 和推测执行如何补充 Lineage 的不足。 下图对比了无 Checkpoint 时的全链重算和有 Checkpoint 时的截断恢复: Lineage 容错原理 每个 RDD 记录了自己的 dependencies——从哪些父 RDD 经过什么变换得来。当一个分区丢失时(Executor 故障、磁盘损坏),Spark 不需要从副本恢复,只需要找到该分区的父 RDD 分区,重新执行 compute 函数。 12345678910RDD-0 (HDFS) → RDD-1 (map) → RDD-2 (filter) → RDD-3 (reduceByK...
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
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...





