Skip to content

大数据技术硬核教程(上篇) ​

认知体系 · 存储底座 · 计算引擎 ​

本篇定位:不是"Hadoop 是什么"的百科式罗列,而是带你穿透每一层技术的设计动机 → 权衡取舍 → 工程落地 → 源码级理解。 实操环境:阿里云 EMR(Hadoop 3.2.1 + Spark 3.5.3 + Hive 3.1.3),所有命令可直接执行。 阅读方式:每章末尾有「🔬 硬核实验」和「💭 思考题」,建议边读边跑。


目录 ​


第 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 / StarRocks

Hadoop 生态(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 ns0.5 秒(心跳一次)
分支预测错误5 ns5 秒
L2 缓存引用7 ns7 秒
互斥锁加锁/解锁25 ns25 秒
主内存引用100 ns1.5 分钟(泡杯咖啡)
内存中压缩 1KB3 μs50 分钟
1Gbps 网络发送 2KB20 μs5.5 小时
NVMe SSD 随机读 4KB~20 μs5.5 小时
内存顺序读 1MB10 μs2.7 小时
同机房网络往返~500 μs6 天
SSD 顺序读 1MB100 μs1.1 天
磁盘寻道(HDD)5 ms2 个月
磁盘顺序读 1MB10 ms4 个月
跨城市网络往返(北京-上海)~30 ms1 年
跨洲网络往返(中国-美西)~150 ms5 年

这张表推导出的核心工程结论:

  1. 顺序访问 >> 随机访问(磁盘上差 100~1000 倍) → 这就是为什么 HDFS 的 Block 是 128MB 而不是 4KB → 这就是为什么 Kafka 敢用磁盘做队列(顺序写) → 这就是为什么 LSM-Tree 打败 B+Tree 成为大数据存储的主流

  2. 网络是新的磁盘 → 这就是为什么有"移动计算而非移动数据"(Data Locality) → 这就是为什么 Shuffle 是 Spark 最大的性能杀手

  3. 内存和磁盘差 10 万倍 → 这就是 Spark 存在的全部理由

约束二:Amdahl 定律 —— 并行化的天花板 ​

加速比 S = 1 / [(1 - p) + p/n]

p = 可并行部分占比
n = 处理器数量

残酷的推论:如果你的程序有 5% 是串行的,哪怕给你无限多的机器,最大加速比也只有 20 倍。

python
# 直观感受一下
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)PCEC宁可不可用也不返回错误元数据
HBasePCEC强一致,单 Region 单点写
CassandraPAEL可调一致性,默认追求低延迟
Kafka (acks=all)PCECISR 机制保证强一致
Kafka (acks=1)PAEL牺牲一致性换吞吐
ZooKeeperPCECZAB 协议,线性一致读需 sync
Redis ClusterPAEL异步复制,可能丢数据

🔬 硬核实验 1.1:亲手制造一次数据丢失

在你的 EMR 集群上(或本地 Docker)搭一个 3 节点 Kafka,设置 acks=1,生产数据的同时 kill -9 leader 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 集群可以这样开启:

bash
hdfs 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
T4AI 融合、架构简化数据治理、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│
└─────────┘            └─────────┘            └─────────┘

关键认知点(面试高频):

  1. NameNode 内存中保存什么?

    • 文件目录树(Namespace)
    • 文件 → Block 列表的映射
    • Block → DataNode 位置的映射(这个不持久化!由 DataNode 汇报重建)
  2. 为什么 Block 位置不持久化? 因为 DataNode 随时可能挂掉/迁移,持久化的位置信息会过期。NameNode 重启后进入安全模式,等待 DataNode 汇报块信息重建 BlockMap。

  3. 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 分布 ​

bash
# 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 次)✅ 能扛整机架故障✅ 较优

验证机架配置:

bash
# 查看当前机架拓扑
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 亿文件

四种解法:

bash
# 解法 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 关键差异(必须知道的坑):

特性HDFSS3/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 列式:一张图讲清本质 ​

假设有这样一张表:

idnameagecity
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) ​

sql
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
         → 可能有匹配,需要读取

这就是为什么数据的物理排序极其重要:

python
# ❌ 未排序:每个 Row Group 的 min/max 都覆盖全值域,剪枝完全失效
df.write.parquet("/data/unsorted")

# ✅ 按过滤列排序:min/max 区间窄,剪枝率大幅提升
df.sort("age").write.parquet("/data/sorted")

🔬 硬核实验 4.1:量化排序带来的性能差异

python
from 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 配置:

python
spark.conf.set("spark.sql.parquet.compression.codec", "zstd")

技术三:Dremel 嵌套编码(Repetition/Definition Levels) ​

这是 Parquet 最难懂但最精妙的部分,来自 Google Dremel 论文。

问题:列式存储怎么表达嵌套结构?比如:

json
{
  "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:怎么选? ​

维度ParquetORC
出身Twitter + ClouderaHortonworks(为 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:三种格式的空间与速度对比 ​

python
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 边界智能切分 Stage

Pipeline 融合示例:

scala
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 字节码 → 执行

亲眼看到优化过程:

python
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)的问题:

java
// 每一行数据,都要走一遍虚函数调用链
while (parent.hasNext()) {
    Row r = parent.next();     // 虚函数调用
    if (filter.eval(r)) {      // 虚函数调用
        out = project.eval(r); // 虚函数调用
    }
}
// 1 亿行数据 = 3 亿次虚函数调用,CPU 分支预测失败率极高

WSCG 的解法:把整个 Stage 的算子编译成一个巨大的 for 循环的 Java 代码,运行时用 Janino 编译成字节码。

java
// 生成的代码大致长这样(无虚函数调用!)
while (scan_hasNext()) {
    long id = scan_getLong(0);
    int age = scan_getInt(1);
    if (age > 90) {                 // filter 内联进来了
        hashAgg_doAggregate(age);   // 聚合也内联
    }
}

效果:CPU 指令数降低一个数量级,接近手写代码性能。

查看生成的代码(硬核操作):

python
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.2Hash ShuffleM × R文件数爆炸(1000 Map × 1000 Reduce = 100 万文件)
Spark 1.2+Sort ShuffleM × 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 调优核心参数 ​

python
# 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 的威力量化测试 ​

python
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 + VeloxIntel + MetaSpark 物理计划下沉到 C++ Velox 引擎执行生产可用,2-3x 加速
PhotonDatabricks闭源 C++ 向量化引擎商业产品,宣称 3-8x
Apache DataFusion CometApple + 社区基于 Rust DataFusion 加速 Spark快速发展中
Blaze快手开源Rust 原生 Spark 执行引擎国内生产验证

核心思想:

JVM 的局限:
  ✗ 无法精细控制内存布局
  ✗ 无法用 SIMD 指令(AVX-512)
  ✗ GC 不可控

C++/Rust 原生引擎:
  ✓ 列式内存布局(Arrow 格式)
  ✓ SIMD 向量化:一条指令处理 8/16 个数据
  ✓ 手动内存管理,零 GC

SIMD 效果直观理解:

标量执行:  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):

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?

维度YARNKubernetes
隔离性cgroup 弱隔离容器强隔离
环境依赖依赖节点上装的 Python/库镜像自包含,环境一致
弹性依赖固定节点池秒级扩缩容,可用 Spot 实例
混部只能跑大数据在离线混部,提升利用率
生态大数据专用统一基础设施

提交 Spark 作业到 K8s:

bash
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.jar

K8s 模式的关键难点: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 项目设定:电商用户行为分析 ​

业务需求:

  1. 计算每日各品类的 PV/UV
  2. 计算用户转化漏斗(浏览 → 加购 → 下单 → 支付)
  3. 找出 Top 100 高价值用户
  4. 输出可供 BI 查询的结果表

技术要求:数据量 1 亿行,要求总耗时 < 10 分钟

7.2 Step 1:生成测试数据(1 亿行) ​

python
# 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()

提交执行:

bash
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

验证:

bash
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 生态) ​

sql
-- 在 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 ​

sql
-- ❌ 朴素写法: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:转化漏斗 ​

sql
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 高价值用户(窗口函数) ​

sql
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(按收益排序) ​

python
# ===== 第 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:结果写出与验证 ​

python
# 写出到 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% → 重点优化它
StagesTask Duration 的 Min/Median/MaxMax >> Median → 数据倾斜
StagesShuffle Read/Write SizeShuffle 量 > 输入量 → 考虑预聚合
StagesSpill (Memory/Disk)有 Disk Spill → 内存不足,增大 executor-memory 或分区数
ExecutorsGC Time / Task TimeGC 占比 >10% → 内存配置有问题
ExecutorsFailed Tasks频繁失败 → 看日志,可能是 OOM
SQL物理计划中的 Exchange 数Exchange 越多越慢
SQLScan 节点的 "number of files read"文件数 >> 预期 → 小文件问题

🔬 硬核实验 7.1:故意制造并诊断三种典型故障 ​

python
# 故障 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 的实时数仓

📮 给读者的话

大数据技术看起来组件繁多,但核心原理就那几条:分而治之、冗余容错、批量摊薄、局部性优先。 所有新框架都是这些原理在不同约束下的重新组合。

当你看到一个新组件时,问自己三个问题:

  1. 它解决的是哪个物理约束?
  2. 它牺牲了什么来换取这个能力?
  3. 我的场景里这个权衡划算吗?

能稳定回答这三个问题,你就从"会用工具的人"变成了"能做架构决策的人"。

基于 Vite 强力驱动 | 纯静态轻量托管