大数据技术硬核教程(上篇)
认知体系 · 存储底座 · 计算引擎
本篇定位:不是"Hadoop 是什么"的百科式罗列,而是带你穿透每一层技术的设计动机 → 权衡取舍 → 工程落地 → 源码级理解。 实操环境:阿里云 EMR(Hadoop 3.2.1 + Spark 3.5.3 + Hive 3.1.3),所有命令可直接执行。 阅读方式:每章末尾有「🔬 硬核实验」和「💭 思考题」,建议边读边跑。
目录
- 第 0 章:先破除三个认知误区
- 第 1 章:大数据技术的第一性原理
- 第 2 章:技术全景图与演进史(2003-2026)
- 第 3 章:分布式存储 —— HDFS 深度解剖
- 第 4 章:列式存储革命 —— Parquet / ORC 位级剖析
- 第 5 章:计算引擎演进 —— 从 MapReduce 到 Spark 到向量化
- 第 6 章:资源调度 —— YARN / K8s 双轨制
- 第 7 章:硬核实战 —— 从 0 构建一个亿级数据分析管道
第 0 章:先破除三个认知误区
在讲任何技术之前,必须先拆掉三堵墙。这三个误区会让你在学习路上走歪至少两年。
❌ 误区一:"大数据 = 数据量大"
真相:数据量只是触发条件,不是本质。
判断你是否需要大数据技术,只有一个标准:
单机的 CPU / 内存 / 磁盘 / 网络,在你可接受的时间和成本内,能否完成这个计算?举个反直觉的例子:
| 场景 | 数据量 | 是否需要大数据技术 | 原因 |
|---|---|---|---|
| 每天 500GB 日志,只做 grep 统计 | 500GB | 不需要 | 一台 128G 内存的机器 + DuckDB,30 分钟跑完 |
| 10GB 数据做全量两两相似度计算 | 10GB | 需要 | 复杂度 O(n²),单机算到天荒地老 |
| 1TB 数据,要求 100ms 内出结果 | 1TB | 需要 | 延迟约束逼出分布式 |
| 200GB 数据,每天跑一次 | 200GB | 大概率不需要 | ClickHouse 单机轻松吃下 |
💡 2024 年以来行业出现了显著的"去大数据化"回潮:
DuckDB、Polars、ClickHouse 单机版的崛起,让大量原本上 Spark 集群的任务回归单机。 有篇著名的论文叫 "Big Data is Dead"(MotherDuck, 2023),核心论点是:大多数公司的大多数查询,处理的数据量在 100GB 以内。给你的启示:学大数据要学"什么时候不用大数据",这是高级工程师和初级工程师的分水岭。
❌ 误区二:"学会 Hadoop 就是学会大数据"
Hadoop 在 2026 年的地位,类似于 C 语言之于现代编程:你必须懂它的原理,但生产环境你很可能不直接写它。
真实的现代数据栈长这样:
存储:S3 / OSS / HDFS (对象存储正在吃掉 HDFS)
表格式:Iceberg / Hudi / Delta Lake (这一层是 2020 年后新增的,极其重要)
计算:Spark / Flink / Trino / DuckDB
调度:Airflow / Dagster / DolphinScheduler
建模:dbt / SQLMesh
OLAP:ClickHouse / Doris / StarRocksHadoop 生态(HDFS + YARN + MapReduce + Hive)只占了其中一小块,而且 MapReduce 已基本退役。
❌ 误区三:"大数据是纯粹的工程活,没什么理论深度"
恰恰相反。大数据技术是分布式系统理论最密集的落地场:
- CAP 定理、PACELC 定理
- 一致性模型(线性一致性、顺序一致性、最终一致性)
- 共识算法(Paxos / Raft / ZAB)
- LSM-Tree 与 B+Tree 的读写放大权衡
- 向量化执行与 CPU 缓存局部性
- 布隆过滤器、HyperLogLog、Count-Min Sketch 等概率数据结构
- 无锁并发、内存屏障、堆外内存管理
本教程会把这些理论掰开揉碎地讲,并且每一个都配可运行的代码验证。
第 1 章:大数据技术的第一性原理
1.1 一切复杂度的源头:三个物理约束
所有大数据技术,本质上都在对抗三个无法绕过的物理事实:
约束一:延迟的数量级鸿沟
这是 Jeff Dean 那张著名的表格(2026 年更新版):
| 操作 | 延迟 | 类比(放大 10 亿倍,人类尺度) |
|---|---|---|
| L1 缓存引用 | 0.5 ns | 0.5 秒(心跳一次) |
| 分支预测错误 | 5 ns | 5 秒 |
| L2 缓存引用 | 7 ns | 7 秒 |
| 互斥锁加锁/解锁 | 25 ns | 25 秒 |
| 主内存引用 | 100 ns | 1.5 分钟(泡杯咖啡) |
| 内存中压缩 1KB | 3 μs | 50 分钟 |
| 1Gbps 网络发送 2KB | 20 μs | 5.5 小时 |
| NVMe SSD 随机读 4KB | ~20 μs | 5.5 小时 |
| 内存顺序读 1MB | 10 μs | 2.7 小时 |
| 同机房网络往返 | ~500 μs | 6 天 |
| SSD 顺序读 1MB | 100 μs | 1.1 天 |
| 磁盘寻道(HDD) | 5 ms | 2 个月 |
| 磁盘顺序读 1MB | 10 ms | 4 个月 |
| 跨城市网络往返(北京-上海) | ~30 ms | 1 年 |
| 跨洲网络往返(中国-美西) | ~150 ms | 5 年 |
这张表推导出的核心工程结论:
顺序访问 >> 随机访问(磁盘上差 100~1000 倍) → 这就是为什么 HDFS 的 Block 是 128MB 而不是 4KB → 这就是为什么 Kafka 敢用磁盘做队列(顺序写) → 这就是为什么 LSM-Tree 打败 B+Tree 成为大数据存储的主流
网络是新的磁盘 → 这就是为什么有"移动计算而非移动数据"(Data Locality) → 这就是为什么 Shuffle 是 Spark 最大的性能杀手
内存和磁盘差 10 万倍 → 这就是 Spark 存在的全部理由
约束二:Amdahl 定律 —— 并行化的天花板
加速比 S = 1 / [(1 - p) + p/n]
p = 可并行部分占比
n = 处理器数量残酷的推论:如果你的程序有 5% 是串行的,哪怕给你无限多的机器,最大加速比也只有 20 倍。
# 直观感受一下
def amdahl(p, n):
return 1 / ((1 - p) + p / n)
for p in [0.5, 0.9, 0.95, 0.99, 0.999]:
print(f"并行比例 {p:>6.1%}: "
f"10核={amdahl(p,10):>6.2f}x "
f"100核={amdahl(p,100):>6.2f}x "
f"1000核={amdahl(p,1000):>7.2f}x "
f"无限核={amdahl(p,10**9):>8.2f}x")输出:
并行比例 50.0%: 10核= 1.82x 100核= 1.98x 1000核= 2.00x 无限核= 2.00x
并行比例 90.0%: 10核= 5.26x 100核= 9.17x 1000核= 9.91x 无限核= 10.00x
并行比例 95.0%: 10核= 6.90x 100核= 16.81x 1000核= 19.63x 无限核= 20.00x
并行比例 99.0%: 10核= 9.17x 100核= 50.25x 1000核= 90.99x 无限核= 100.00x
并行比例 99.9%: 10核= 9.91x 100核= 90.99x 1000核= 500.25x 无限核= 1000.00x工程启示:Spark 作业跑得慢,加机器往往没用。要先找出串行瓶颈(通常是:Driver 端单点收集、数据倾斜导致的长尾 Task、写出阶段的单文件合并)。
约束三:CAP 与 PACELC —— 分布式的原罪
CAP 定理(Brewer, 2000):一致性(C)、可用性(A)、分区容错性(P),三者最多同时满足两个。
但 CAP 经常被误读。更准确的表述是 PACELC(Abadi, 2012):
if (Partition 分区发生) {
在 Availability 和 Consistency 之间选择 → PA / PC
} else {
在 Latency 和 Consistency 之间选择 → EL / EC
}主流系统的 PACELC 定位:
| 系统 | 分区时 | 正常时 | 说明 |
|---|---|---|---|
| HDFS NameNode (HA) | PC | EC | 宁可不可用也不返回错误元数据 |
| HBase | PC | EC | 强一致,单 Region 单点写 |
| Cassandra | PA | EL | 可调一致性,默认追求低延迟 |
| Kafka (acks=all) | PC | EC | ISR 机制保证强一致 |
| Kafka (acks=1) | PA | EL | 牺牲一致性换吞吐 |
| ZooKeeper | PC | EC | ZAB 协议,线性一致读需 sync |
| Redis Cluster | PA | EL | 异步复制,可能丢数据 |
🔬 硬核实验 1.1:亲手制造一次数据丢失
在你的 EMR 集群上(或本地 Docker)搭一个 3 节点 Kafka,设置
acks=1,生产数据的同时kill -9leader broker,观察数据丢失。然后改成acks=all+min.insync.replicas=2重做一遍。这个实验会让你对"一致性"这个词有肌肉记忆级的理解。
1.2 大数据技术的三大设计范式
所有大数据系统,都在反复使用这三招:
范式一:分而治之(Divide & Conquer)
大问题 → 切分成独立小问题 → 并行求解 → 合并结果- MapReduce:Map 切分 → Reduce 合并
- HDFS:文件切成 Block 分散存储
- 分区表:按时间/哈希切分数据
关键难点:怎么切?切不均就是数据倾斜(本教程中篇会花整章讲)。
范式二:冗余换可靠(Redundancy for Reliability)
单机 MTTF(平均无故障时间)= 3 年
1000 台机器的集群 → 平均每天挂 1 台结论:在大规模集群里,故障是常态而非异常。
应对手段:
- 副本(Replication):HDFS 3 副本,存储成本 300%
- 纠删码(Erasure Coding):RS(6,3) 编码,存储成本仅 150%,但恢复时 CPU/网络开销大
📐 纠删码数学原理:把数据切成 6 个数据块,通过 Reed-Solomon 编码生成 3 个校验块。任意丢失 3 个块都能恢复。 本质是解线性方程组:在 GF(2^8) 有限域上,用范德蒙德矩阵构造编码矩阵,丢失块 = 解方程的未知数。
HDFS 3.x 已原生支持 EC,适合冷数据。你的 EMR 集群可以这样开启:
bashhdfs ec -listPolicies hdfs ec -enablePolicy -policy RS-6-3-1024k hdfs ec -setPolicy -path /cold_data -policy RS-6-3-1024k
范式三:批量摊薄(Batching & Amortization)
单次操作成本高?那就攒一批一起做。
- HDFS:128MB Block 摊薄寻址成本
- Kafka:Producer 端 batch.size 攒批发送
- LSM-Tree:内存 MemTable 攒够了再刷盘
- 向量化执行:一次处理 1024 行而非 1 行(摊薄函数调用和虚函数开销)
- GPU 计算:SIMT 模型,一次调度一个 warp(32 线程)
💭 思考题 1.1:批量能提升吞吐,但会损害什么?请从 Kafka 的
linger.ms参数思考「吞吐 vs 延迟」的本质权衡。
第 2 章:技术全景图与演进史(2003-2026)
2.1 四个时代的分水岭
┌─────────────────────────────────────────────────────────────────┐
│ T1 2003-2012 蛮荒时代:Google 三驾马车 & Hadoop 复刻 │
│ GFS(2003) → MapReduce(2004) → BigTable(2006) │
│ Hadoop(2006) / HBase(2008) / Hive(2009) / Kafka(2011) │
│ 关键词:能跑起来就行,磁盘密集,批处理为王 │
├─────────────────────────────────────────────────────────────────┤
│ T2 2012-2018 内存时代:Spark 统治 & 流批分裂 │
│ Spark(2012) / YARN(2012) / Flink(2014) / Presto(2013) │
│ Lambda 架构(批+流双链路)盛行 │
│ 关键词:内存计算,SQL 化,Lambda 架构的复杂度之痛 │
├─────────────────────────────────────────────────────────────────┤
│ T3 2018-2023 云原生 & 湖仓时代 │
│ Delta Lake(2019) / Iceberg(2018) / Hudi(2019) │
│ 存算分离 / K8s 化 / Snowflake & Databricks 崛起 │
│ ClickHouse / Doris / StarRocks 实时 OLAP 爆发 │
│ 关键词:湖仓一体,存算分离,ACID 进入数据湖 │
├─────────────────────────────────────────────────────────────────┤
│ T4 2023-2026 AI 原生 & 单机复兴 │
│ 向量数据库 / RAG 数据管道 / Feature Store │
│ DuckDB / Polars 单机复兴 / Arrow 统一内存格式 │
│ Iceberg 成为事实标准(AWS/Snowflake/Databricks 全面拥抱) │
│ Text-to-SQL / AI Agent 直接操作数据栈 │
│ 关键词:AI 与数据栈融合,Arrow 生态统一,简化回归 │
└─────────────────────────────────────────────────────────────────┘2.2 每个时代解决了什么、留下了什么问题
| 时代 | 解决的核心问题 | 留下的新问题 | 催生的下一代技术 |
|---|---|---|---|
| T1 | 廉价机器存算 PB 级数据 | 慢(磁盘 IO)、难写(MR 编程模型反人类) | Spark、Hive |
| T2 | 快(内存)、易写(SQL/DataFrame) | 流批两套代码、数据湖无 ACID、小文件 | Flink 统一、Iceberg |
| T3 | 湖上 ACID、存算分离弹性 | 组件过多、成本失控、延迟仍不够低 | 湖仓融合、实时 OLAP |
| T4 | AI 融合、架构简化 | 数据治理、AI 幻觉、成本归因 | 待续(你的时代) |
2.3 2026 年的现代数据栈全景
┌──────────────────────────────────────────────────────────────────┐
│ 应用层 / 消费层 │
│ BI(Superset/Tableau) │ AI Agent │ RAG │ 推荐 │ 报表 │ 实时大屏 │
└────────────────────────────┬─────────────────────────────────────┘
│
┌────────────────────────────▼─────────────────────────────────────┐
│ 查询 / 服务层 │
│ Trino/Presto │ ClickHouse │ Doris │ StarRocks │ DuckDB │
│ 向量库: Milvus / Qdrant / pgvector │
└────────────────────────────┬─────────────────────────────────────┘
│
┌────────────────────────────▼─────────────────────────────────────┐
│ 计算 / 处理层 │
│ 批: Spark │ 流: Flink │ 单机: DuckDB/Polars │
│ 建模: dbt / SQLMesh 编排: Airflow / Dagster / DolphinScheduler │
└────────────────────────────┬─────────────────────────────────────┘
│
┌────────────────────────────▼─────────────────────────────────────┐
│ ★ 表格式层(Table Format)★ ← 这层是关键! │
│ Apache Iceberg │ Apache Hudi │ Delta Lake │
│ 提供:ACID / Schema 演进 / 时间旅行 / 分区演进 / 行级更新 │
├──────────────────────────────────────────────────────────────────┤
│ 文件格式层(File Format) │
│ Parquet(列式,OLAP 主流) │ ORC │ Avro(行式,流式友好) │ Arrow │
├──────────────────────────────────────────────────────────────────┤
│ 存储层(Storage) │
│ 对象存储: S3 / OSS / MinIO │ 分布式文件: HDFS │ 本地 NVMe │
└──────────────────────────────────────────────────────────────────┘
▲ ▲
│ │
┌───────┴──────────┐ ┌──────────┴────────────┐
│ 数据接入层 │ │ 元数据 & 治理 │
│ Kafka/Pulsar │ │ Hive Metastore │
│ Flink CDC │ │ AWS Glue / Nessie │
│ DataX/Seatunnel │ │ DataHub / OpenMetadata│
└──────────────────┘ └───────────────────────┘🎯 给学习者的路线建议: 不要按图从下往上学。推荐路径是:
SQL → Hive/Spark SQL → HDFS/Parquet 原理 → Spark 内核 → Iceberg → Flink → OLAP 引擎 → 治理与成本即:先能干活,再懂原理,再懂架构。
第 3 章:分布式存储 —— HDFS 深度解剖
3.1 HDFS 的设计哲学:五个"不"
HDFS 之所以长成这样,源于它旗帜鲜明的取舍:
| 设计取舍 | 原因 | 代价 |
|---|---|---|
| 不支持随机写(只能追加) | 简化一致性模型,避免分布式锁 | 更新数据必须重写整个文件 |
| 不适合小文件 | 每个文件的元数据约占 NameNode 150 字节内存 | 千万级小文件会撑爆 NameNode |
| 不追求低延迟 | 为高吞吐优化,128MB 大 Block | 不能做在线服务的存储 |
| 不做 POSIX 完全兼容 | 放弃兼容性换性能 | 不能直接 mount 当普通文件系统用 |
| 不假设硬件可靠 | 廉价机器,故障是常态 | 3 副本导致 200% 存储浪费 |
核心假设一句话总结:
"一次写入,多次读取(Write-Once-Read-Many),大文件顺序扫描,硬件会坏。"
3.2 架构解剖:三个角色的真实职责
┌──────────────────────────┐
│ NameNode │
│ ┌────────────────────┐ │
元数据操作 │ │ FsImage (全量快照) │ │
┌──────────────────► │ │ EditLog (增量日志) │ │
│ │ │ BlockMap (内存) │◄─┼── 心跳 + 块报告
│ │ └────────────────────┘ │
│ └──────────────────────────┘
┌───┴────┐ │
│ Client │ │ 返回 DataNode 地址列表
└───┬────┘ ◄───────────────────────┘
│
│ ★ 真实数据传输不经过 NameNode ★
│
▼
┌─────────┐ pipeline ┌─────────┐ pipeline ┌─────────┐
│DataNode1│ ─────────► │DataNode2│ ─────────► │DataNode3│
└─────────┘ └─────────┘ └─────────┘关键认知点(面试高频):
NameNode 内存中保存什么?
- 文件目录树(Namespace)
- 文件 → Block 列表的映射
- Block → DataNode 位置的映射(这个不持久化!由 DataNode 汇报重建)
为什么 Block 位置不持久化? 因为 DataNode 随时可能挂掉/迁移,持久化的位置信息会过期。NameNode 重启后进入安全模式,等待 DataNode 汇报块信息重建 BlockMap。
FsImage 和 EditLog 的关系
FsImage = 某个时刻的完整元数据快照(类似数据库的全量备份) EditLog = 之后所有的修改操作日志(类似 binlog) 重启恢复 = FsImage + 重放 EditLog问题:EditLog 越来越大,重启越来越慢 → 需要定期 Checkpoint 合并
- Hadoop 1.x:SecondaryNameNode 做合并
- Hadoop 2.x+ HA 模式:StandbyNameNode 做合并(更优雅)
3.3 写流程深度追踪:Pipeline 与 Packet
这是 HDFS 最精妙的设计之一。
Client 写 1 个 128MB Block 的完整流程:
1. Client → NameNode: create() 请求
2. NameNode: 检查权限、租约(Lease),返回 3 个 DataNode 地址
↓
3. Client 建立 Pipeline:Client → DN1 → DN2 → DN3
↓
4. 数据切分:
Block(128MB) → 多个 Packet(64KB) → 每个 Packet 含多个 Chunk(512B) + Checksum(4B)
↓
5. 流水线传输:
Client 发 Packet1 给 DN1
DN1 收到后【立即转发】给 DN2,同时写本地磁盘 ← 流水线的精髓
DN2 收到后【立即转发】给 DN3,同时写本地磁盘
↓
6. ACK 反向流回:DN3 → DN2 → DN1 → Client
↓
7. 全部 Packet 完成 → Client → NameNode: complete()为什么用 Pipeline 而不是 Client 同时写 3 份?
方案 A(Client 并发写 3 份):
Client 上行带宽 = 3 × 数据量 ← Client 网卡成为瓶颈
方案 B(Pipeline 串行转发):
每个节点上行带宽 = 1 × 数据量 ← 带宽均摊,充分利用集群网络这是典型的把负载从单点分散到链路的设计思想。
🔬 硬核实验 3.1:观察真实的 Block 分布
# 1. 生成一个 300MB 的测试文件(会被切成 3 个 Block)
dd if=/dev/urandom of=/tmp/big.dat bs=1M count=300
# 2. 上传到 HDFS
hdfs dfs -mkdir -p /lab/hdfs
hdfs dfs -put /tmp/big.dat /lab/hdfs/
# 3. 查看 Block 分布(关键命令!)
hdfs fsck /lab/hdfs/big.dat -files -blocks -locations
# 你会看到类似输出:
# /lab/hdfs/big.dat 314572800 bytes, replicated: replication=2, 3 block(s)
# 0. BP-xxx:blk_1073741825_1001 len=134217728 Live_repl=2
# [DatanodeInfoWithStorage[172.16.146.221:9866,DISK],
# DatanodeInfoWithStorage[172.16.146.222:9866,DISK]]
# 1. BP-xxx:blk_1073741826_1002 len=134217728 ...
# 2. BP-xxx:blk_1073741827_1003 len=46137344 ... ← 最后一个 Block 不足 128MB
# 4. 在 DataNode 上找到真实的物理文件
find /mnt/disk1/hdfs -name "blk_1073741825*"
# 会看到两个文件:
# blk_1073741825 ← 真实数据
# blk_1073741825_1001.meta ← 校验和文件
# 5. 验证:Block 就是普通文件,可以直接拼接还原
cat blk_1073741825 blk_1073741826 blk_1073741827 > /tmp/restored.dat
md5sum /tmp/big.dat /tmp/restored.dat # 应该完全一致💡 这个实验的价值:让你亲眼看到「HDFS 只是一个把大文件切块分散存储的系统」,破除神秘感。
3.4 副本放置策略:机架感知的数学
HDFS 默认的 3 副本放置策略(Rack Awareness):
副本 1:Client 所在节点(若 Client 在集群外,随机选一个)
副本 2:与副本 1 【不同机架】 的某个节点
副本 3:与副本 2 【同机架】 的另一个节点为什么是这个组合?一个精妙的三方权衡:
| 策略 | 写带宽消耗 | 机架故障容忍 | 读取本地性 |
|---|---|---|---|
| 3 副本全同机架 | 低(跨机架 0 次) | ❌ 机架挂=数据全丢 | ✅ 优 |
| 3 副本全不同机架 | 高(跨机架 2 次) | ✅ 优 | ❌ 差 |
| HDFS 策略(1+2 分布) | 中(跨机架 1 次) | ✅ 能扛整机架故障 | ✅ 较优 |
验证机架配置:
# 查看当前机架拓扑
hdfs dfsadmin -printTopology
# 输出示例:
# Rack: /default-rack
# 172.16.146.221:9866 (core-1-1)
# 172.16.146.222:9866 (core-1-2)⚠️ 单机架集群(如你的 EMR 默认配置)副本策略会退化。生产环境的多机架集群必须配置
net.topology.script.file.name指向机架映射脚本。
3.5 HDFS 的三大痛点与工业界解法
痛点一:小文件问题(最经典的坑)
数学计算:
NameNode 单条元数据 ≈ 150 字节(文件 inode + block 信息)
1 亿个小文件 ≈ 1亿 × 150B ≈ 15 GB 内存
而 NameNode 是 JVM 进程,堆内存超过 32GB 后 GC 停顿严重(指针压缩失效)
→ 实际上单 NameNode 管理上限约 1~3 亿文件四种解法:
# 解法 1:HAR 归档(Hadoop Archive)
hadoop archive -archiveName logs.har -p /input/logs /output
# 解法 2:合并小文件(Spark 层)
spark.sql("SET spark.sql.adaptive.enabled=true")
spark.sql("SET spark.sql.adaptive.coalescePartitions.enabled=true")
# 或显式 repartition
df.repartition(10).write.parquet("/output")
# 解法 3:使用 Iceberg / Hudi 的自动小文件合并(现代最优解)
# Iceberg:
CALL catalog.system.rewrite_data_files(
table => 'db.table',
options => map('target-file-size-bytes','134217728')
);
# 解法 4:Federation(多 NameNode 分治)—— 治标不治本,运维复杂痛点二:NameNode 单点
HA 架构(生产必备):
┌──────────────┐ ┌──────────────┐
│ Active NN │ │ Standby NN │
└──────┬───────┘ └──────┬───────┘
│ 写 EditLog │ 读 EditLog
▼ ▼
┌──────────────────────────────────────┐
│ JournalNode 集群 (通常 3 或 5 个) │
│ 基于 Quorum 机制,多数派写入成功即可 │
└──────────────────────────────────────┘
▲ ▲
│ │
┌──────┴───────────────────────┴───────┐
│ ZooKeeper (ZKFC 故障检测与选主) │
└──────────────────────────────────────┘脑裂(Split-Brain)防护:JournalNode 使用 Epoch Number(纪元号)机制——新 Active NN 上任时递增 epoch,JournalNode 拒绝所有低 epoch 的写请求。这是 Paxos/Raft 中 term 概念的同构实现。
痛点三:存算耦合,弹性差
2020 年后的答案:存算分离
传统 HDFS: 云原生存算分离:
┌──────────────┐ ┌──────────────┐
│ 计算 + 存储 │ │ 计算(弹性) │ ← 按需扩缩容,用完释放
│ 绑定在一起 │ └──────┬───────┘
│ │ │ S3/OSS API
│ 扩容=同时扩 │ ┌──────▼───────┐
│ 缩容=数据迁移 │ │ 对象存储(廉价) │ ← 独立扩展,成本 1/3
└──────────────┘ └──────────────┘对象存储 vs HDFS 关键差异(必须知道的坑):
| 特性 | HDFS | S3/OSS |
|---|---|---|
| rename 操作 | O(1),元数据操作 | O(n),实际是复制+删除! |
| 一致性 | 强一致 | 现已强一致(2020 年 S3 升级后) |
| list 性能 | 快 | 慢,且有 API 调用费用 |
| 随机读 | 支持 | 支持(Range 请求)但有额外延迟 |
⚠️ 著名的 S3 rename 陷阱:Spark 写数据时的
FileOutputCommitter v1依赖 rename 做原子提交,在 S3 上会变成巨慢的复制操作。 解法:使用S3A Committer(magic committer)或直接上 Iceberg(它用元数据文件原子切换,天然免疫)。
💭 思考题 3.1:为什么 HDFS 的 Block 默认从 64MB 变成了 128MB,而没有变成 1GB? 提示:从「寻址开销占比」和「Map Task 并行度」两个角度分析,找出那个甜点区间。
第 4 章:列式存储革命 —— Parquet / ORC 位级剖析
4.1 行式 vs 列式:一张图讲清本质
假设有这样一张表:
| id | name | age | city |
|---|---|---|---|
| 1 | 张三 | 25 | 北京 |
| 2 | 李四 | 30 | 上海 |
| 3 | 王五 | 28 | 北京 |
行式存储(磁盘上的字节排列):
[1,张三,25,北京][2,李四,30,上海][3,王五,28,北京]
└────一行────┘ └────一行────┘ └────一行────┘列式存储:
[1,2,3][张三,李四,王五][25,30,28][北京,上海,北京]
└id列┘ └───name列───┘ └─age列─┘ └───city列────┘执行 SELECT AVG(age) FROM t 时:
| 行式 | 列式 | |
|---|---|---|
| 需读取的数据 | 全部 4 列 | 仅 age 列(1/4 的 IO) |
| CPU 缓存友好度 | 差(数据类型混杂) | 优(同类型连续,可 SIMD) |
| 压缩率 | 低(类型混杂,模式少) | 高(同列同类型,模式重复) |
4.2 Parquet 文件的物理结构(逐层拆解)
┌─────────────────────────────────────────────────┐
│ PAR1 ← 4 字节魔数 │
├─────────────────────────────────────────────────┤
│ Row Group 0 (默认 128MB,对齐 HDFS Block) │
│ ┌───────────────────────────────────────────┐ │
│ │ Column Chunk: id │ │
│ │ ┌─────────────────────────────────────┐ │ │
│ │ │ Page 0 (默认 1MB) │ │ │
│ │ │ ├ PageHeader (含统计信息 min/max) │ │ │
│ │ │ ├ Repetition Levels (嵌套结构用) │ │ │
│ │ │ ├ Definition Levels (NULL 处理用) │ │ │
│ │ │ └ Encoded Values (编码后的值) │ │ │
│ │ ├─────────────────────────────────────┤ │ │
│ │ │ Page 1 ... │ │ │
│ │ └─────────────────────────────────────┘ │ │
│ ├───────────────────────────────────────────┤ │
│ │ Column Chunk: name │ │
│ ├───────────────────────────────────────────┤ │
│ │ Column Chunk: age │ │
│ └───────────────────────────────────────────┘ │
├─────────────────────────────────────────────────┤
│ Row Group 1 ... │
├─────────────────────────────────────────────────┤
│ ★ Footer(元数据,读文件时先读这里)★ │
│ - Schema 定义 │
│ - 每个 Row Group / Column Chunk 的: │
│ offset、压缩前后大小、编码方式 │
│ ★ min/max/null_count 统计信息 ★ │
├─────────────────────────────────────────────────┤
│ Footer Length (4 字节) │
│ PAR1 ← 结尾魔数 │
└─────────────────────────────────────────────────┘为什么元数据在文件末尾(Footer)而不是开头?
因为写入时是流式的——写完所有数据才知道每个 Chunk 的实际大小和统计信息。 读取时:先 seek 到文件末尾 → 读 4 字节 Footer Length → 反向 seek 读 Footer → 得到全部元数据 → 精准 seek 到需要的 Column Chunk。 两次 seek 换来极致的按需读取。
4.3 三大核心优化技术
技术一:谓词下推 + 统计信息剪枝(Predicate Pushdown)
SELECT * FROM users WHERE age > 60执行时 Parquet Reader 会:
读 Footer → 发现 Row Group 0 的 age 列 min=18, max=35
→ 35 < 60,整个 Row Group 直接跳过,一个字节都不读!
→ Row Group 1 的 age 列 min=20, max=80
→ 可能有匹配,需要读取这就是为什么数据的物理排序极其重要:
# ❌ 未排序:每个 Row Group 的 min/max 都覆盖全值域,剪枝完全失效
df.write.parquet("/data/unsorted")
# ✅ 按过滤列排序:min/max 区间窄,剪枝率大幅提升
df.sort("age").write.parquet("/data/sorted")🔬 硬核实验 4.1:量化排序带来的性能差异
pythonfrom pyspark.sql import SparkSession from pyspark.sql.functions import rand, floor, col import time spark = SparkSession.builder.appName("parquet-lab").getOrCreate() # 造 5000 万行数据 df = (spark.range(0, 50_000_000) .withColumn("age", floor(rand(42) * 100)) .withColumn("score", rand(7) * 1000) .withColumn("city", floor(rand(1) * 300))) # A: 无序写入 df.write.mode("overwrite").parquet("hdfs:///lab/parquet/unsorted") # B: 排序写入 df.sort("age").write.mode("overwrite").parquet("hdfs:///lab/parquet/sorted") def bench(path, label): t = time.time() n = spark.read.parquet(path).filter(col("age") > 95).count() print(f"{label}: {n} 行, 耗时 {time.time()-t:.2f}s") bench("hdfs:///lab/parquet/unsorted", "无序") bench("hdfs:///lab/parquet/sorted", "排序")典型结果:排序版本通常快 3~10 倍。用 Spark UI 的 SQL 页面查看 "number of files read" 和 "scan bytes" 指标,能直接看到剪枝效果。
技术二:编码压缩(Encoding)—— 这是 Parquet 的黑魔法
Parquet 在通用压缩(Snappy/ZSTD)之前,先做编码,效果远超单纯压缩:
| 编码方式 | 原理 | 适用场景 | 压缩效果 |
|---|---|---|---|
| Dictionary(字典) | 建立值→ID 字典,存 ID | 低基数列(如城市、性别) | 极佳,10-100x |
| RLE(游程编码) | 连续相同值存 (值, 次数) | 排序后的列、稀疏列 | 极佳 |
| Bit-Packing(位打包) | 只用必要的 bit 数存整数 | 小范围整数 | 好 |
| Delta Encoding(差分) | 存相邻值的差 | 时间戳、递增 ID | 极佳 |
| Byte Stream Split | 浮点数按字节位拆分重组 | Float/Double 列 | 中等提升 |
实例计算:
一个存储 1 亿条记录的 "city" 列(共 300 个不同城市)
朴素存储:1亿 × 平均 9 字节(UTF-8中文) = 900 MB
字典编码:300 个城市字典(2.7KB) + 1亿 × 9 bit(表示0-299) ≈ 112 MB
再 RLE(若已排序):≈ 几 KB
最后 ZSTD 压缩:进一步减半
最终:900MB → 可能不到 10MB💡 实用建议:写 Parquet 时优先选
zstd(比 snappy 压缩率高 20-30%,速度接近),Spark 配置:pythonspark.conf.set("spark.sql.parquet.compression.codec", "zstd")
技术三:Dremel 嵌套编码(Repetition/Definition Levels)
这是 Parquet 最难懂但最精妙的部分,来自 Google Dremel 论文。
问题:列式存储怎么表达嵌套结构?比如:
{
"name": "张三",
"phones": ["138xxx", "139xxx"],
"address": { "city": "北京", "zip": null }
}Parquet 的解法:用两个"元数据数字"配合值列表,无损还原嵌套结构。
Definition Level (D):这个值"定义"到了嵌套的第几层(处理 NULL)
Repetition Level (R):这个值在第几层"重复"(处理数组)以 phones 列为例:
值: "138xxx" "139xxx"
R: 0 1 ← 0 表示新记录开始,1 表示同一记录内的重复
D: 2 2 ← 定义到第 2 层(非 null)为什么这么设计? 因为它让每一列都能完全独立地读取和重建,不需要读其他列——这正是列式存储的核心诉求。
📚 延伸阅读:Google Dremel 论文《Dremel: Interactive Analysis of Web-Scale Datasets》(VLDB 2010),是列式存储领域最重要的论文之一。
4.4 Parquet vs ORC:怎么选?
| 维度 | Parquet | ORC |
|---|---|---|
| 出身 | Twitter + Cloudera | Hortonworks(为 Hive 优化) |
| 生态支持 | 更广(Spark/Flink/Trino/Python/Rust 全支持) | Hive/Spark 好,其他一般 |
| 压缩率 | 好 | 略好(轻量级索引更精细) |
| 索引 | Row Group 级 min/max + 可选 Bloom Filter | 三级索引(File/Stripe/Row Group 10000行) |
| ACID 支持 | 靠上层(Iceberg/Delta) | Hive ACID 原生支持更好 |
| 云原生 | 事实标准 | 相对弱 |
| 2026 现状 | 推荐默认选择 | Hive 老栈继续用 |
一句话结论:新项目无脑选 Parquet;已有 Hive ORC 表没必要迁移。
🔬 硬核实验 4.2:三种格式的空间与速度对比
from pyspark.sql.functions import rand, floor, col, concat, lit
import subprocess, time
df = (spark.range(0, 20_000_000)
.withColumn("age", floor(rand(1)*100))
.withColumn("city", concat(lit("city_"), floor(rand(2)*300)))
.withColumn("amount", rand(3)*10000))
formats = {
"csv": lambda p: df.write.mode("overwrite").csv(p),
"json": lambda p: df.write.mode("overwrite").json(p),
"parquet-snappy": lambda p: df.write.mode("overwrite")
.option("compression","snappy").parquet(p),
"parquet-zstd": lambda p: df.write.mode("overwrite")
.option("compression","zstd").parquet(p),
"orc": lambda p: df.write.mode("overwrite").orc(p),
}
for name, writer in formats.items():
path = f"hdfs:///lab/fmt/{name}"
t = time.time()
writer(path)
wt = time.time() - t
size = subprocess.check_output(
f"hdfs dfs -du -s -h {path}", shell=True).decode().split()[0:2]
print(f"{name:18s} 写入 {wt:6.1f}s 体积 {' '.join(size)}")典型结果(2000 万行):
csv 写入 45.2s 体积 1.1 G
json 写入 78.6s 体积 3.4 G
parquet-snappy 写入 28.3s 体积 156 M
parquet-zstd 写入 31.7s 体积 112 M
orc 写入 30.1s 体积 108 M💥 冲击性结论:JSON 比 Parquet 大 30 倍。这就是为什么「日志直接存 JSON 上大数据平台」是新手最烧钱的错误之一。
第 5 章:计算引擎演进 —— 从 MapReduce 到 Spark 到向量化
5.1 MapReduce:伟大但已退役
MapReduce 的历史功绩:它第一次让普通工程师能写分布式程序,不用关心容错、调度、数据分发。
但它有三个致命缺陷:
缺陷一:每个阶段都落盘
MapReduce 的一次迭代:
读HDFS → Map → 写本地磁盘 → Shuffle(网络) → Reduce → 写HDFS
↑落盘 ↑落盘
做 10 次迭代(如机器学习):
10 × (读HDFS + 写HDFS) = 20 次磁盘往返Spark 的改进:中间结果保留在内存 → 迭代计算快 10~100 倍。
缺陷二:编程模型太原始
实现一个简单的 SELECT ... GROUP BY ... JOIN ...,MapReduce 需要写几百行 Java,串联多个 Job。
缺陷三:JVM 进程启动开销
每个 Task 启动一个 JVM,启动开销约 1~3 秒。小任务的开销占比可能超过 50%。
📌 2026 年的定位:MapReduce 只剩下"理解分布式计算模型"的教学价值。生产环境请直接用 Spark/Flink。 但 Shuffle 的思想是 MapReduce 留给世界的最大遗产,Spark/Flink 的 Shuffle 都是它的演进版。
5.2 Spark 核心:为什么快?(四个层次的答案)
层次一:内存计算(最表层的答案,也是最常被误解的)
⚠️ 辟谣:Spark 不是"把所有数据加载到内存"。Spark 默认也是流式处理分区的,内存放不下会溢写磁盘(spill)。 真正的差异是:Stage 之间的中间结果不强制落 HDFS。
层次二:DAG 调度与惰性求值
MapReduce: 每个 Job 是孤立的,无法全局优化
Spark: 构建完整 DAG 后统一优化
→ 能做 Pipeline 融合(多个 map 合并成一个 Task 执行)
→ 能做谓词下推、列裁剪
→ 能在 Shuffle 边界智能切分 StagePipeline 融合示例:
rdd.map(f).filter(g).map(h)
// 不会产生 3 个 Task 阶段,而是融合成一个:
// for each record: h(g_filter(f(record)))
// 数据只遍历一次,无中间集合创建层次三:Catalyst 优化器(Spark SQL 的灵魂)
SQL / DataFrame 代码
↓
Unresolved Logical Plan (解析语法,列名还未绑定)
↓ ← Catalog 元数据
Resolved Logical Plan (列名/表名解析完成)
↓ ← 规则优化 (RBO)
Optimized Logical Plan (谓词下推、常量折叠、列裁剪、子查询消除)
↓ ← 代价优化 (CBO) + 统计信息
Physical Plans (多个候选) (选 BroadcastHashJoin 还是 SortMergeJoin?)
↓ ← 代价模型选择最优
Selected Physical Plan
↓ ← 全阶段代码生成 (WSCG)
生成 Java 字节码 → 执行亲眼看到优化过程:
df = spark.read.parquet("hdfs:///lab/parquet/sorted")
result = df.filter(col("age") > 90).select("id", "age").groupBy("age").count()
# 查看完整的四阶段计划
result.explain(True)
# 只看物理计划(日常调优最常用)
result.explain("formatted")你会在输出里看到关键信息:
== Physical Plan ==
*(2) HashAggregate(keys=[age], functions=[count(1)])
+- Exchange hashpartitioning(age, 200) ← ★ Shuffle 发生在这
+- *(1) HashAggregate(keys=[age], functions=[partial_count(1)]) ← ★ Map 端预聚合
+- *(1) Filter (isnotnull(age) AND (age > 90))
+- FileScan parquet [id,age] ← ★ 只读 2 列(列裁剪生效)
PushedFilters: [IsNotNull(age), GreaterThan(age,90)] ← ★ 谓词下推生效🎯 读懂物理计划是 Spark 调优的第一技能。记住三个关键字:
Exchange= Shuffle(性能杀手,越少越好)PushedFilters= 谓词下推(有则好)*(n)前缀 = WholeStageCodegen 生效(有则好)
层次四:Tungsten 引擎 —— 榨干 CPU 与内存
这是最深的一层,也是 Spark 2.0 之后真正的性能飞跃来源。
三大技术:
1️⃣ 堆外内存 + 自定义二进制格式(UnsafeRow)
传统 JVM 对象存储 "abc" 这个字符串:
对象头(16B) + hash(4B) + char[]引用(8B) + 数组对象头(16B) + 6字节数据
= 约 48+ 字节,其中真实数据只占 6 字节 → 开销 800%
Tungsten UnsafeRow(紧凑二进制):
[null bitmap][固定长度字段][变长字段偏移+长度][变长数据]
→ 接近 C 语言结构体的紧凑度,且 GC 完全不管(堆外)收益:内存占用降低 50-70%,GC 停顿大幅减少。
2️⃣ Whole-Stage Code Generation(全阶段代码生成)
传统火山模型(Volcano Model)的问题:
// 每一行数据,都要走一遍虚函数调用链
while (parent.hasNext()) {
Row r = parent.next(); // 虚函数调用
if (filter.eval(r)) { // 虚函数调用
out = project.eval(r); // 虚函数调用
}
}
// 1 亿行数据 = 3 亿次虚函数调用,CPU 分支预测失败率极高WSCG 的解法:把整个 Stage 的算子编译成一个巨大的 for 循环的 Java 代码,运行时用 Janino 编译成字节码。
// 生成的代码大致长这样(无虚函数调用!)
while (scan_hasNext()) {
long id = scan_getLong(0);
int age = scan_getInt(1);
if (age > 90) { // filter 内联进来了
hashAgg_doAggregate(age); // 聚合也内联
}
}效果:CPU 指令数降低一个数量级,接近手写代码性能。
查看生成的代码(硬核操作):
result.explain("codegen") # 会打印出生成的 Java 源码,非常震撼3️⃣ 缓存感知计算(Cache-aware Computation)
排序时把 key 和指针打包成 8 字节放在一起,排序时只比较 key 前缀,减少随机内存访问 → 大幅提升 CPU L1/L2 缓存命中率。
5.3 Shuffle 深度剖析:Spark 最大的性能瓶颈
Shuffle 为什么慢?
Shuffle = 磁盘 IO + 网络 IO + 序列化/反序列化 + 排序
一次 Shuffle 涉及:
Map 端:计算分区 → 排序 → 写本地磁盘(可能多次 spill + merge)
Reduce 端:网络拉取 → 反序列化 → 归并排序 → 计算Shuffle 的演进史
| 版本 | 机制 | 小文件数 | 问题 |
|---|---|---|---|
| Spark <1.2 | Hash Shuffle | M × R | 文件数爆炸(1000 Map × 1000 Reduce = 100 万文件) |
| Spark 1.2+ | Sort Shuffle | M × 2 | 每个 Map 输出 1 个数据文件 + 1 个索引文件 |
| Spark 3.2+ | Push-Based Shuffle | 更优 | Map 端主动推送,Reduce 端顺序读,减少随机 IO |
| 外部方案 | Remote Shuffle Service(Celeborn/Uniffle) | 存算分离 | K8s 弹性场景必备 |
🔥 2026 前沿:Apache Celeborn(原 RSS,阿里开源)已成为 Spark on K8s 的标配。它把 Shuffle 数据存到独立服务,使 Executor 可以随时被回收而不丢 Shuffle 数据,实现真正的弹性。
Shuffle 调优核心参数
# 1. 分区数(最重要的参数)
spark.conf.set("spark.sql.shuffle.partitions", "200") # 默认 200,需按数据量调整
# 经验公式:分区数 ≈ 总数据量 / 128MB,且应为 Executor 总核数的 2-3 倍
# 2. 开启 AQE(Spark 3.x 必开!)
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true") # 自动合并小分区
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true") # 自动处理倾斜
# 3. Shuffle 读写缓冲
spark.conf.set("spark.shuffle.file.buffer", "64k") # 默认 32k,调大减少磁盘写次数
spark.conf.set("spark.reducer.maxSizeInFlight", "96m") # 默认 48m
# 4. 序列化(Kryo 比 Java 序列化快 10 倍)
spark.conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")🔬 硬核实验 5.1:AQE 的威力量化测试
from pyspark.sql.functions import rand, when, col, floor, lit
import time
# 构造一个严重倾斜的数据集:80% 的数据 key 都是 0
skewed = (spark.range(0, 30_000_000)
.withColumn("key", when(rand(1) < 0.8, lit(0))
.otherwise(floor(rand(2) * 1000)))
.withColumn("val", rand(3)))
dim = spark.range(0, 1000).withColumnRenamed("id", "key") \
.withColumn("name", col("key").cast("string"))
def run(aqe_enabled, label):
spark.conf.set("spark.sql.adaptive.enabled", str(aqe_enabled).lower())
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", str(aqe_enabled).lower())
# 禁用广播,强制走 SortMergeJoin 以暴露倾斜
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
t = time.time()
skewed.join(dim, "key").groupBy("name").count().count()
print(f"{label}: {time.time()-t:.1f}s")
run(False, "AQE 关闭")
run(True, "AQE 开启")打开 Spark UI 的 Stages 页面,对比两次运行的 Task Duration 分布图——AQE 关闭时你会看到一根长得离谱的柱子(倾斜的那个 Task),开启后会被自动拆分。
5.4 前沿:向量化执行与原生引擎(2024-2026 最大变化)
即便有了 Tungsten,Spark 仍受制于 JVM。于是出现了用 C++ 重写执行层的浪潮:
| 项目 | 主导方 | 原理 | 状态 |
|---|---|---|---|
| Gluten + Velox | Intel + Meta | Spark 物理计划下沉到 C++ Velox 引擎执行 | 生产可用,2-3x 加速 |
| Photon | Databricks | 闭源 C++ 向量化引擎 | 商业产品,宣称 3-8x |
| Apache DataFusion Comet | Apple + 社区 | 基于 Rust DataFusion 加速 Spark | 快速发展中 |
| Blaze | 快手开源 | Rust 原生 Spark 执行引擎 | 国内生产验证 |
核心思想:
JVM 的局限:
✗ 无法精细控制内存布局
✗ 无法用 SIMD 指令(AVX-512)
✗ GC 不可控
C++/Rust 原生引擎:
✓ 列式内存布局(Arrow 格式)
✓ SIMD 向量化:一条指令处理 8/16 个数据
✓ 手动内存管理,零 GCSIMD 效果直观理解:
标量执行: for (i=0; i<8; i++) c[i] = a[i] + b[i]; → 8 条指令
SIMD 执行: _mm256_add_epi32(a, b); → 1 条指令💡 给学习者的信号:如果你想在 2026 年后做大数据内核方向,Rust + Arrow + DataFusion 是最值得投入的技术栈组合。
第 6 章:资源调度 —— YARN / K8s 双轨制
6.1 YARN 架构本质
┌───────────────────────────────────────────────────┐
│ ResourceManager (RM) —— 全局资源大管家 │
│ ├─ Scheduler:只负责分配资源,不关心任务逻辑 │
│ └─ ApplicationsManager:管理所有 App 的生命周期 │
└──────────────┬────────────────────────────────────┘
│
┌───────────┼───────────┐
▼ ▼ ▼
┌──────┐ ┌──────┐ ┌──────┐
│ NM │ │ NM │ │ NM │ NodeManager:单机资源管家
│ │ │ │ │ │
│┌────┐│ │┌────┐│ │┌────┐│
││ AM ││ ││Ctnr││ ││Ctnr││ AM = ApplicationMaster(每个 App 一个)
│└────┘│ │└────┘│ │└────┘│ Ctnr = Container(资源容器)
└──────┘ └──────┘ └──────┘YARN 的核心设计突破:把 资源管理 和 任务调度 解耦。
- RM 只管"给谁多少 CPU 内存"
- 每个应用自己的 AM 管"这些资源怎么用"
这使得 YARN 能同时跑 Spark、Flink、MapReduce、Tez 等不同计算框架。
6.2 三种调度器对比
| 调度器 | 策略 | 适用场景 |
|---|---|---|
| FIFO | 先进先出 | 测试环境,生产禁用 |
| Capacity Scheduler | 队列预分配容量,队列内 FIFO | 多部门共享集群(主流) |
| Fair Scheduler | 动态公平分配,所有 App 均分资源 | 交互式查询多的场景 |
生产队列配置示例(capacity-scheduler.xml):
<property>
<name>yarn.scheduler.capacity.root.queues</name>
<value>prod,dev,adhoc</value>
</property>
<property>
<name>yarn.scheduler.capacity.root.prod.capacity</name>
<value>60</value> <!-- 保底 60% -->
</property>
<property>
<name>yarn.scheduler.capacity.root.prod.maximum-capacity</name>
<value>90</value> <!-- 空闲时最多借到 90% -->
</property>
<property>
<name>yarn.scheduler.capacity.root.adhoc.capacity</name>
<value>15</value>
</property>
<property>
<name>yarn.scheduler.capacity.root.adhoc.maximum-am-resource-percent</name>
<value>0.1</value> <!-- 限制 AM 占比,防止大量小任务占满 -->
</property>6.3 Spark on K8s:2026 年的主流方向
为什么要从 YARN 迁到 K8s?
| 维度 | YARN | Kubernetes |
|---|---|---|
| 隔离性 | cgroup 弱隔离 | 容器强隔离 |
| 环境依赖 | 依赖节点上装的 Python/库 | 镜像自包含,环境一致 |
| 弹性 | 依赖固定节点池 | 秒级扩缩容,可用 Spot 实例 |
| 混部 | 只能跑大数据 | 在离线混部,提升利用率 |
| 生态 | 大数据专用 | 统一基础设施 |
提交 Spark 作业到 K8s:
spark-submit \
--master k8s://https://<k8s-api-server>:6443 \
--deploy-mode cluster \
--name spark-etl \
--class com.example.ETLJob \
--conf spark.executor.instances=10 \
--conf spark.kubernetes.container.image=registry.cn-hangzhou.aliyuncs.com/xxx/spark:3.5.3 \
--conf spark.kubernetes.namespace=data \
--conf spark.kubernetes.authenticate.driver.serviceAccountName=spark \
--conf spark.kubernetes.executor.deleteOnTermination=true \
--conf spark.dynamicAllocation.enabled=true \
--conf spark.dynamicAllocation.shuffleTracking.enabled=true \
local:///opt/jobs/etl.jarK8s 模式的关键难点:Shuffle 数据问题
Executor 被回收 → 它本地的 Shuffle 数据丢失 → 下游 Task 失败重算
解法:
方案 1:spark.dynamicAllocation.shuffleTracking.enabled=true(保守,不回收有 shuffle 数据的 Executor)
方案 2:★ Remote Shuffle Service(Celeborn/Uniffle)← 生产推荐💭 思考题 6.1:你的 EMR 集群是 YARN 模式。假设公司要迁到 K8s,你会列出哪 5 个必须先解决的问题?(提示:日志采集、Shuffle、镜像管理、权限体系、成本归因)
第 7 章:硬核实战 —— 从 0 构建一个亿级数据分析管道
现在把前面所有知识串起来,在你的 EMR 集群上做一个完整的项目。
7.1 项目设定:电商用户行为分析
业务需求:
- 计算每日各品类的 PV/UV
- 计算用户转化漏斗(浏览 → 加购 → 下单 → 支付)
- 找出 Top 100 高价值用户
- 输出可供 BI 查询的结果表
技术要求:数据量 1 亿行,要求总耗时 < 10 分钟
7.2 Step 1:生成测试数据(1 亿行)
# gen_data.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import (
col, rand, floor, expr, lit, when, date_add,
to_date, concat, unix_timestamp, from_unixtime
)
spark = (SparkSession.builder
.appName("gen-ecommerce-data")
.config("spark.sql.shuffle.partitions", "200")
.config("spark.sql.parquet.compression.codec", "zstd")
.enableHiveSupport()
.getOrCreate())
N = 100_000_000 # 1 亿行
df = (spark.range(0, N)
# 用户 ID:500 万用户,故意做成幂律分布(模拟真实的长尾)
.withColumn("user_id",
when(rand(1) < 0.3, floor(rand(2) * 10000)) # 30% 流量来自 1 万活跃用户
.otherwise(floor(rand(3) * 5_000_000)))
# 商品 ID:100 万商品
.withColumn("item_id", floor(rand(4) * 1_000_000))
# 品类:20 个品类
.withColumn("category_id", floor(rand(5) * 20))
# 行为类型
.withColumn("behavior",
when(rand(6) < 0.70, lit("pv"))
.when(rand(6) < 0.88, lit("cart"))
.when(rand(6) < 0.96, lit("order"))
.otherwise(lit("pay")))
# 金额(仅 pay 行为有效)
.withColumn("amount",
when(col("behavior") == "pay", (rand(7) * 2000 + 10).cast("decimal(10,2)"))
.otherwise(lit(None).cast("decimal(10,2)")))
# 时间:最近 30 天
.withColumn("ts", from_unixtime(
unix_timestamp(lit("2026-08-20 00:00:00")) + floor(rand(8) * 30 * 86400)))
.withColumn("dt", to_date(col("ts")))
.drop("id"))
# ★ 关键:按分区列写出,并控制每个分区的文件数
(df.repartition(200, col("dt")) # 按天重分区,避免小文件
.write
.mode("overwrite")
.partitionBy("dt") # 分区表
.parquet("hdfs:///lab/ecommerce/user_behavior"))
print("数据生成完成")
spark.stop()提交执行:
spark-submit \
--master yarn --deploy-mode cluster \
--driver-memory 2g \
--executor-memory 6g \
--executor-cores 3 \
--num-executors 4 \
--conf spark.sql.adaptive.enabled=true \
--conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
gen_data.py验证:
hdfs dfs -du -s -h /lab/ecommerce/user_behavior
hdfs dfs -ls /lab/ecommerce/user_behavior | head
# 应该看到 dt=2026-08-20 / dt=2026-08-21 ... 这样的分区目录7.3 Step 2:注册为 Hive 表(打通 Hive 生态)
-- 在 spark-sql 或 beeline 中执行
CREATE DATABASE IF NOT EXISTS ecom;
CREATE EXTERNAL TABLE IF NOT EXISTS ecom.user_behavior (
user_id BIGINT,
item_id BIGINT,
category_id BIGINT,
behavior STRING,
amount DECIMAL(10,2),
ts STRING
)
PARTITIONED BY (dt DATE)
STORED AS PARQUET
LOCATION 'hdfs:///lab/ecommerce/user_behavior';
-- ★ 关键:修复分区元数据(否则查不到数据,这是超高频的坑)
MSCK REPAIR TABLE ecom.user_behavior;
-- 验证
SELECT dt, COUNT(*) FROM ecom.user_behavior GROUP BY dt ORDER BY dt LIMIT 5;⚠️ 经典踩坑点:外部表建好后不执行
MSCK REPAIR TABLE,查询会返回 0 行。因为 Hive Metastore 不知道有哪些分区。 Spark 也可以用:spark.sql("ALTER TABLE ecom.user_behavior RECOVER PARTITIONS")
7.4 Step 3:核心分析 SQL(含性能讲解)
需求 1:每日各品类 PV/UV
-- ❌ 朴素写法:COUNT(DISTINCT) 会导致严重的单点聚合
SELECT dt, category_id,
COUNT(*) AS pv,
COUNT(DISTINCT user_id) AS uv
FROM ecom.user_behavior
WHERE behavior = 'pv' AND dt >= '2026-08-20'
GROUP BY dt, category_id;
-- ✅ 优化写法 1:两阶段去重(先按 user 去重再计数)
SELECT dt, category_id,
SUM(cnt) AS pv,
COUNT(*) AS uv
FROM (
SELECT dt, category_id, user_id, COUNT(*) AS cnt
FROM ecom.user_behavior
WHERE behavior = 'pv' AND dt >= '2026-08-20'
GROUP BY dt, category_id, user_id
) t
GROUP BY dt, category_id;
-- ✅ 优化写法 2:近似去重(误差 <2%,速度提升 10 倍以上)
SELECT dt, category_id,
COUNT(*) AS pv,
APPROX_COUNT_DISTINCT(user_id, 0.02) AS uv_approx
FROM ecom.user_behavior
WHERE behavior = 'pv' AND dt >= '2026-08-20'
GROUP BY dt, category_id;📐 APPROX_COUNT_DISTINCT 背后的算法:HyperLogLog
核心洞察:随机哈希值的前导零个数能估计基数。
- 如果你观察到最大前导零数为 k,说明大约看过 2^k 个不同元素
- 用多个桶(bucket)分别统计再调和平均,降低方差
- 空间复杂度 O(log log n):统计 10 亿基数只需约 12KB 内存!
这是概率数据结构的经典应用,同类的还有:
- Bloom Filter:判断"一定不存在"(用于 Join 预过滤、HBase 读优化)
- Count-Min Sketch:估计元素频次(用于 TopK、热点检测)
- T-Digest:估计分位数(用于 P99 延迟计算)
需求 2:转化漏斗
WITH user_stage AS (
SELECT
dt,
user_id,
MAX(CASE WHEN behavior = 'pv' THEN 1 ELSE 0 END) AS has_pv,
MAX(CASE WHEN behavior = 'cart' THEN 1 ELSE 0 END) AS has_cart,
MAX(CASE WHEN behavior = 'order' THEN 1 ELSE 0 END) AS has_order,
MAX(CASE WHEN behavior = 'pay' THEN 1 ELSE 0 END) AS has_pay
FROM ecom.user_behavior
WHERE dt >= '2026-08-20'
GROUP BY dt, user_id
)
SELECT
dt,
SUM(has_pv) AS pv_users,
SUM(has_cart) AS cart_users,
SUM(has_order) AS order_users,
SUM(has_pay) AS pay_users,
ROUND(SUM(has_cart) / NULLIF(SUM(has_pv),0), 4) AS pv_to_cart,
ROUND(SUM(has_order) / NULLIF(SUM(has_cart),0), 4) AS cart_to_order,
ROUND(SUM(has_pay) / NULLIF(SUM(has_order),0),4) AS order_to_pay
FROM user_stage
GROUP BY dt
ORDER BY dt;性能要点:这个 SQL 只有一次 Shuffle(GROUP BY dt, user_id),外层聚合可以在同分区完成。用 EXPLAIN 验证只有 2 个 Exchange。
需求 3:Top 100 高价值用户(窗口函数)
WITH user_metrics AS (
SELECT
user_id,
SUM(CASE WHEN behavior='pay' THEN amount ELSE 0 END) AS total_amount,
COUNT(CASE WHEN behavior='pay' THEN 1 END) AS pay_cnt,
COUNT(DISTINCT dt) AS active_days,
COUNT(DISTINCT category_id) AS category_breadth
FROM ecom.user_behavior
WHERE dt >= '2026-08-20'
GROUP BY user_id
)
SELECT *,
-- 综合价值分(简易 RFM 变体)
ROUND(total_amount * 0.5 + pay_cnt * 100 * 0.3 + active_days * 50 * 0.2, 2) AS value_score,
RANK() OVER (ORDER BY total_amount DESC) AS amount_rank
FROM user_metrics
WHERE total_amount > 0
ORDER BY value_score DESC
LIMIT 100;⚠️ 窗口函数的性能陷阱:
OVER (ORDER BY x)不带PARTITION BY会把所有数据 Shuffle 到一个分区做全局排序,是典型的单点瓶颈。 优化思路:先用ORDER BY ... LIMIT做 TopN(Spark 有专门的 TakeOrderedAndProject 算子,只在各分区取 TopN 再合并),再对小结果集做排名。
7.5 Step 4:性能调优实操 —— 把 10 分钟优化到 3 分钟
调优 Checklist(按收益排序)
# ===== 第 1 优先级:减少读取的数据量 =====
# ① 分区裁剪:WHERE 条件必须包含分区列
# ✅ WHERE dt >= '2026-08-20' → 只扫描对应分区
# ❌ WHERE ts >= '2026-08-20' → 全表扫描!
# ② 列裁剪:只 select 需要的列(Parquet 自动生效)
# ❌ SELECT * → 读所有列
# ✅ SELECT a, b → 只读 2 列
# ===== 第 2 优先级:减少 Shuffle =====
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.localShuffleReader.enabled", "true")
# 广播小表(消除 Shuffle Join)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "50m") # 默认 10m
# ===== 第 3 优先级:并行度匹配资源 =====
# 公式:shuffle.partitions ≈ executor数 × executor核数 × 2~3
# 也应满足:单分区数据量 ≈ 128~256MB
total_cores = 4 * 3 # num-executors × executor-cores
spark.conf.set("spark.sql.shuffle.partitions", str(total_cores * 3))
# ===== 第 4 优先级:内存与序列化 =====
spark.conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
spark.conf.set("spark.memory.fraction", "0.7") # 执行+存储内存占比
spark.conf.set("spark.sql.files.maxPartitionBytes", "134217728") # 单分区读取上限
# ===== 第 5 优先级:缓存复用数据 =====
df_base = spark.table("ecom.user_behavior").filter("dt >= '2026-08-20'")
df_base.cache() # 多次使用时才 cache,且要注意内存是否够
df_base.count() # 触发缓存物化
# ... 多次使用 df_base ...
df_base.unpersist() # 用完释放资源配置的黄金法则
给定集群:2 个 core 节点 × (4 vCPU, 16 GB) ← 你的 EMR 配置
可用资源估算(YARN 会预留一部分给系统):
总核数 ≈ 2 × 4 = 8,实际可用约 6~7
总内存 ≈ 2 × 16 = 32 GB,YARN 可分配约 24 GB
配置建议:
--executor-cores 3 ← 3~5 最优(太大 HDFS 并发差,太小并行度低)
--num-executors 3 ← (可用核数 - Driver占用) / executor-cores
--executor-memory 5g ← (可用内存 / executor数) × 0.9,留出 overhead
--conf spark.executor.memoryOverhead=1g ← 堆外内存,约为 executor-memory 的 10~20%
--driver-memory 2g🎯 为什么 executor-cores 推荐 3~5?
- 太大(如 16):单个 Executor 内多线程争抢 HDFS 客户端,吞吐下降;且 GC 压力集中
- 太小(如 1):无法共享广播变量和缓存,内存浪费;且 JVM 数量多,开销大
- 3~5 是社区多年实践得出的经验区间
7.6 Step 5:结果写出与验证
# 写出到 Hive 结果表
result.write \
.mode("overwrite") \
.format("parquet") \
.option("compression", "zstd") \
.saveAsTable("ecom.dws_category_daily")
# ★ 写出后必做:检查小文件
# hdfs dfs -count /user/hive/warehouse/ecom.db/dws_category_daily
# 如果文件数 >> 分区数,说明有小文件问题
# 优化写出(控制文件数与大小)
result.repartition(10) \
.write.mode("overwrite") \
.saveAsTable("ecom.dws_category_daily")
# 或使用 Spark 3.x 的自动优化
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128m")7.7 全流程监控:在 Spark UI 里该看什么
| 页面 | 关键指标 | 异常信号 |
|---|---|---|
| Jobs | 每个 Job 耗时分布 | 某个 Job 占总时间 >50% → 重点优化它 |
| Stages | Task Duration 的 Min/Median/Max | Max >> Median → 数据倾斜 |
| Stages | Shuffle Read/Write Size | Shuffle 量 > 输入量 → 考虑预聚合 |
| Stages | Spill (Memory/Disk) | 有 Disk Spill → 内存不足,增大 executor-memory 或分区数 |
| Executors | GC Time / Task Time | GC 占比 >10% → 内存配置有问题 |
| Executors | Failed Tasks | 频繁失败 → 看日志,可能是 OOM |
| SQL | 物理计划中的 Exchange 数 | Exchange 越多越慢 |
| SQL | Scan 节点的 "number of files read" | 文件数 >> 预期 → 小文件问题 |
🔬 硬核实验 7.1:故意制造并诊断三种典型故障
# 故障 1:数据倾斜
# 构造 90% 数据同一个 key,跑 groupBy,在 Stages 页面看 Task 时长分布
# 故障 2:OOM
# executor-memory 设成 512m,跑一个大 collect(),看报错栈
# 故障 3:小文件
# spark.sql.shuffle.partitions 设成 2000,写出后 hdfs dfs -count 看文件数
# 每种故障都:① 复现 ② 在 UI 里找到症状 ③ 修复 ④ 验证
# 这个练习的价值远超读十篇调优文章📋 上篇总结:你现在应该掌握的能力清单
理论层
- [ ] 能说出大数据技术的三个物理约束(延迟鸿沟、Amdahl、CAP/PACELC)
- [ ] 能解释为什么 HDFS 的 Block 是 128MB
- [ ] 能画出 HDFS 写流程的 Pipeline 并说明为什么这样设计
- [ ] 能解释列式存储为什么压缩率高、扫描快(三个原因)
- [ ] 能说清 Catalyst 的四个优化阶段
- [ ] 能解释 Tungsten 的三大技术及其解决的问题
- [ ] 知道 HyperLogLog / Bloom Filter 的原理和适用场景
实操层
- [ ] 能用
hdfs fsck查看 Block 分布并找到物理文件 - [ ] 能读懂
explain("formatted")的物理计划,识别 Exchange / PushedFilters - [ ] 能根据集群规格计算出合理的 executor 配置
- [ ] 能在 Spark UI 中通过 Task 时长分布诊断数据倾斜
- [ ] 能独立完成"建表 → 分区修复 → 分析 → 写出 → 调优"全流程
- [ ] 知道
MSCK REPAIR TABLE这类高频坑
判断力层
- [ ] 能判断一个需求该不该用大数据技术
- [ ] 能在 Parquet / ORC / CSV / JSON 之间做出有依据的选择
- [ ] 知道 2026 年哪些技术在上升期、哪些在退役
🔮 中篇预告
《大数据技术硬核教程(中篇):实时计算 · 湖仓一体 · 数据建模》
- Flink 深度剖析:Checkpoint 与 Chandy-Lamport 算法、Exactly-Once 的真相、Watermark 与乱序处理、状态后端选型(RocksDB 调优)
- 流批一体的三次尝试:Lambda → Kappa → 湖仓流批一体,每次失败在哪
- Iceberg / Hudi / Delta Lake 三选一:元数据结构逐层拆解、MoR vs CoW、隐藏分区、时间旅行的实现原理
- CDC 全链路:Debezium / Flink CDC 实战,从 MySQL binlog 到湖仓的完整数据流
- 数仓建模:维度建模 vs Data Vault vs One Big Table,分层规范(ODS/DWD/DWS/ADS)的工程落地
- 实时 OLAP 对决:ClickHouse / Doris / StarRocks 架构对比与选型决策树
- 完整实战:MySQL → Flink CDC → Kafka → Iceberg → StarRocks 的实时数仓
📮 给读者的话
大数据技术看起来组件繁多,但核心原理就那几条:分而治之、冗余容错、批量摊薄、局部性优先。 所有新框架都是这些原理在不同约束下的重新组合。
当你看到一个新组件时,问自己三个问题:
- 它解决的是哪个物理约束?
- 它牺牲了什么来换取这个能力?
- 我的场景里这个权衡划算吗?
能稳定回答这三个问题,你就从"会用工具的人"变成了"能做架构决策的人"。