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
内存管理与 Tungsten:堆外内存、序列化与代码生成
上一篇解决了 Shuffle 的物理机制。这一篇进入内存管理——Spark 如何突破 JVM 的内存瓶颈。 Spark 是一个 JVM 应用,天然受制于 Java 的内存模型:对象头开销大、GC 停顿不可控、序列化效率低。Project Tungsten 是 Spark 为解决这三个问题发起的底层优化计划,它从三条线同时推进——堆外内存管理、二进制数据格式、全阶段代码生成。 本文只抓一个问题:Tungsten 的三条优化线分别解决了什么问题,以及 Spark 的统一内存管理模型如何在执行内存和存储内存之间做动态调配。 下图展示了 Executor 的统一内存管理模型——Execution Memory 和 Storage Memory 之间的动态借用机制: Java 对象的内存开销 一个 Java 字符串 “abcd” 在堆内占多少字节? 12345678910111213java.lang.String 对象: 对象头: 12 bytes (64-bit JVM, 压缩指针) hash: 4 bytes (int) value[]: 4 bytes ...
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...
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
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
导读:为什么 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
宽依赖与窄依赖: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...
Announcement
人生只是,守株待兔









