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
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
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] 全部完成 ...
2017-12-07
为什么要自建实时计算平台
#为什么要自建一个离线平台# 可以优化资源利用率。 业务平台应该把精力放在业务上。 #什么是实时计算# 强调响应时间短(相对于离线计算):毫秒级、亚秒级、秒级。T+1 的报表都是离线计算。 数据的价值随着时间的流逝而迅速降低。 常见技术方案: 流计算 + 实时存储 or 消息队列 流计算 + 实现 OLAP #什么是流式计算# 实时且无界。 数据驱动计算,事件触发。 有状态及持续集成。 流计算引擎:Spark Streaming、Flink Streaming、Storm/JStorm、Samza 等。 #Spark Streaming 模型# Micro-Batch 模式。看起来是流式处理的,实际上还是一小批一小批处理的。从批处理走到流处理。 最小延时:batch 的处理时间 最大延时:batch interval(通常2s-10s) + batch 处理时间。 使用场景:数据清洗(实时数据通道)、数据 ETL 等。 对于熟悉 Spark 批处理的 RD 非常容易上手。 #Flink Streaming# Native Streaming。 低延时,通常在毫秒...
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 作...
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...
2025-07-29
经典面试问题的大数据解法——Spark 与 Flink 实战
“100 亿个数中找出最大的 1000 个”、“两个 10GB 的文件找出共同的 URL”——这些经典面试题的本质都是内存放不下。单机方案围绕分治展开,分布式方案则把分治思想映射到集群节点上。本文按问题类型组织,每类问题给出从单机到 Spark/Flink 的渐进式解法,并附上概率数据结构(布隆过滤器、HyperLogLog、Count-Min Sketch)在近似场景中的应用。 引言:大数据问题的共同特征 为什么"内存放不下" 面试中给出的数据规模往往是精心设计的——刚好跨过单机内存的边界: 数据规模 内存需求 典型服务器内存 能否放入内存 1 亿个 int 400 MB 16 GB ✅ 10 亿个 int 4 GB 16 GB ✅(但留给程序的余量不多) 100 亿个 int 40 GB 16 GB ❌ 10 亿个 URL(平均 100 字节) 100 GB 16 GB ❌ 上表只计算了裸数据大小。实际使用 HashMap、HashSet 等容器时,对象头、指针、负载因子会使内存占用膨胀 3-5 倍。 通用解题框架 123...





