Apache Spark 是目前最主流的分布式批处理引擎,其统一的编程模型和内存计算能力使其在大规模数据处理中占据核心地位。本文系统讲解 Spark 三大核心抽象、查询优化器、Shuffle 机制与生产调参策略。
1. Spark 核心抽象对比:RDD vs DataFrame vs Dataset
1.1 三大 API 全景对比
| 特性 | RDD | DataFrame | Dataset |
|---|---|---|---|
| 引入版本 | Spark 1.0 | Spark 1.3 | Spark 1.6 |
| 类型安全 | 是(编译期) | 否(运行时检查) | 是(编译期) |
| 模式推断 | 无 | 有 | 有 |
| 性能优化 | 无内置优化 | Catalyst + Tungsten | Catalyst + Tungsten |
| API 风格 | 函数式 | DSL/SQL | 函数式 + DSL |
| 序列化 | Java 序列化 | Tungsten 二进制 | Tungsten 编码器 |
| 适用语言 | Scala/Java/Python/R | Scala/Java/Python/R | Scala/Java |
选型建议:
- DataFrame(首选):结构化数据批处理,最大化利用 Catalyst 优化器
- Dataset:需要在编译期类型安全且性能优先的 Scala/Java 场景
- RDD:非结构化数据、自定义分区、跨版本兼容或 fine-grained 控制
1.2 RDD(弹性分布式数据集)
// 创建 RDD
val rdd = sparkContext.parallelize(Seq(1, 2, 3, 4, 5))
// 转换(Transformation - 懒执行)
val mappedRDD = rdd.map(x => x * 2)
val filteredRDD = rdd.filter(x => x > 2)
// 行动(Action - 触发执行)
val result = filteredRDD.reduce(_ + _)
// 持久化到内存
mappedRDD.cache()
// 键值对 RDD 操作
val pairRDD = rdd.map(x => (x % 2, x))
val grouped = pairRDD.groupByKey()
val reduced = pairRDD.reduceByKey(_ + _)
Lineage 血缘机制:RDD 通过依赖关系记录计算逻辑,当分区数据丢失时可自动重算,无需全量复制。
// DAG 依赖示例
val textFile = sc.textFile("hdfs://logs/*.log")
val errors = textFile.filter(_.contains("ERROR")) // Narrow dependency
val mapped = errors.map(_.split("\\t")) // Narrow dependency
val reduced = mapped.map(x => (x(0), 1))
.reduceByKey(_ + _) // Wide dependency (Shuffle)
1.3 DataFrame
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
spark = SparkSession.builder.appName("BatchPipeline").getOrCreate()
# 读取 Parquet
df = spark.read.parquet("s3://data/orders/")
# DSL 查询
df.filter(col("amount") > 100) \
.groupBy("category") \
.agg(sum("amount").alias("gmv"), count("*").alias("order_count")) \
.orderBy(desc("gmv")) \
.show()
# SQL 查询(通过 Spark SQL)
df.createOrReplaceTempView("orders")
spark.sql("""
SELECT category, SUM(amount) as gmv, COUNT(*) as cnt
FROM orders
WHERE amount > 100
GROUP BY category
ORDER BY gmv DESC
""").show()
# 保存结果
df.write.mode("overwrite").partitionBy("dt").parquet("s3://output/orders_summary/")
1.4 Dataset(Scala)
case class Order(orderId: Long, userId: Long, amount: Double, category: String, dt: String)
// 编译期类型安全
val ds: Dataset[Order] = spark.read.parquet("s3://data/orders/").as[Order]
// 编译时可检查字段名
ds.filter(_.amount > 100)
.groupByKey(_.category)
.mapGroups { case (cat, iter) =>
val list = iter.toSeq
(cat, list.map(_.amount).sum, list.size)
}
.toDF("category", "gmv", "order_count")
.orderBy(desc("gmv"))
2. Spark SQL 与 Catalyst 优化器
2.1 Catalyst 优化流程
SQL / DataFrame DSL
↓
Unresolved Logical Plan (解析表名、列名)
↓
Analyzer → Resolved Logical Plan
↓
Catalyst Optimizer → Optimized Logical Plan
- 谓词下推 (Predicate Pushdown)
- 列裁剪 (Column Pruning)
- 常量折叠 (Constant Folding)
- 连接重排序 (Join Reordering)
↓
Spark Planner → Physical Plans
↓
Cost Model → Best Physical Plan
↓
Tungsten → Optimized Java Code Generation
↓
RDD Execution
2.2 常用优化规则
| 优化规则 | 说明 | 效果 |
|---|---|---|
| 谓词下推 | 将 WHERE 条件下推到数据源 | 减少读取数据量 |
| 列裁剪 | 只读取查询需要的列 | 减少 I/O |
| 常量折叠 | 编译期计算常量表达式 | 减少运行时计算 |
| 聚合下推 | 部分聚合在 map 端完成 | 减少 Shuffle 数据 |
| Broadcast Join | 小表广播到大表所在节点 | 避免 Hash Shuffle |
| Sort Merge Join | 有序数据直接归并 | 降低内存压力 |
# 查看执行计划
spark.sql("SELECT * FROM orders WHERE amount > 100").explain(extended=True)
# 输出:
# == Parsed Logical Plan ==
# == Analyzed Logical Plan ==
# == Optimized Logical Plan ==
# == Physical Plan ==
# 强制执行 Broadcast Join
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "100MB")
from pyspark.sql.functions import broadcast
joined = largeDF.join(broadcast(smallDF), "user_id")
2.3 Tungsten 执行引擎
Tungsten 通过以下方式提升执行效率:
┌─────────────────────────────────────────────────┐
│ Tungsten 优化 │
├─────────────────────────────────────────────────┤
│ 1. 紧凑二进制格式 (UnsafeRow) → 减少 GC │
│ 2. 代码生成 (Whole-Stage Code Generation) │
│ 3. off-heap 内存管理 → 避免 JVM GC 影响 │
└─────────────────────────────────────────────────┘
// 开启 Whole-Stage Codegen(默认开启)
spark.conf.set("spark.sql.codegen.wholeStage", "true")
// 查看是否使用了 codegen
// Physical Plan 中会出现 *
// *HashAggregate -> 表示该算子参与了 Whole-Stage Codegen
3. Shuffle 机制原理与优化
3.1 Shuffle 过程详解
Shuffle 是 Spark 中消耗最大的操作,涉及磁盘 I/O、网络传输和序列化。
Map Stage Shuffle Stage Reduce Stage
┌─────────┐ ┌─────────────┐
│ MapTask │── write ──→ ┌──────────────┐ ── read →│ ReduceTask │
│ │ shuffle │ Shuffle File │ │ │
│ │ partition │ (磁盘/Memory)│ │ │
│ │── write ──→ │ │ ── read →│ │
└─────────┘ └──────────────┘ └─────────────┘
排序分区 → 溢写到磁盘 → 合并文件 → 网络拉取 → 合并排序
触发 Shuffle 的算子:groupByKey、reduceByKey、aggregateByKey、sortByKey、join、cogroup、repartition、distinct。
3.2 Shuffle 优化策略
| 策略 | 方法 | 适用场景 |
|---|---|---|
| 减少 Shuffle 次数 | 使用 reduceByKey 替代 groupByKey + map | 聚合操作 |
| Map 端预聚合 | aggregateByKey、combineByKey | 键值对聚合 |
| Broadcast Join | 小表广播,避免 Shuffle Join | 大小表关联 |
| Range Partitioner | 倾斜键分区优化 | 数据倾斜 |
| Shuffle 分区数调整 | spark.sql.shuffle.partitions | 任务并行度 |
| 序列化优化 | Kryo 序列化 | 大量对象传输 |
# 坏写法:先 group 再聚合,没有 map 端预聚合
rdd.map(lambda x: (x[0], x[1])).groupByKey().mapValues(sum) # 全部数据 Shuffle
# 好写法:reduceByKey 自带 map 端 combine
rdd.map(lambda x: (x[0], x[1])).reduceByKey(lambda a, b: a + b) # 先 combine 再 Shuffle
# aggregateByKey 更灵活
rdd.aggregateByKey(
zeroValue=(0, 0),
seqFunc=lambda acc, v: (acc[0] + v, acc[1] + 1), # map 端
combFunc=lambda a, b: (a[0] + b[0], a[1] + b[1]) # reduce 端
)
3.3 数据倾斜处理
from pyspark.sql.functions import rand, lit
# 方法一:加盐打散倾斜键
spark.sql("""
SELECT
CONCAT(user_id, '_', CAST(rand() * 10 AS INT)) as salted_key,
amount
FROM orders
""").groupBy("salted_key").agg(sum("amount"))
# 方法二:两阶段聚合
rdd.map(lambda x: (x[0] % 10, x)).groupByKey()... # 先随机局部聚合
rdd.map(lambda x: (x[0], x[1])).reduceByKey() # 再全局聚合
# 方法三:Spark SQL AQE 自动优化
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
4. Spark Core 调参实战
4.1 核心参数配置
# SparkSession 调参示例
spark = SparkSession.builder \
.appName("ProductionBatchJob") \
.master("yarn") \
.config("spark.executor.instances", "50") \
.config("spark.executor.cores", "4") \
.config("spark.executor.memory", "16g") \
.config("spark.executor.memoryOverhead", "4g") \
.config("spark.driver.memory", "8g") \
.config("spark.sql.shuffle.partitions", "400") \
.config("spark.default.parallelism", "200") \
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \
.config("spark.sql.autoBroadcastJoinThreshold", "100MB") \
.config("spark.sql.files.maxPartitionBytes", "128MB") \
.getOrCreate()
4.2 参数速查表
| 参数名 | 推荐值 | 说明 |
|---|---|---|
spark.executor.instances | 集群核数/executor.cores | 总 Executor 数量 |
spark.executor.cores | 2-5 | 每个 Executor 的 CPU 核数 |
spark.executor.memory | 8-32g | 堆内存大小 |
spark.executor.memoryOverhead | memory * 0.1-0.25 | off-heap/Native/Netty 内存 |
spark.sql.shuffle.partitions | 2-4倍 executor 总数 | Shuffle 后分区数 |
spark.default.parallelism | 2-3倍 executor 总数 | RDD 默认分区数 |
spark.sql.adaptive.enabled | true | 自适应查询执行 (AQE) |
spark.serializer | KryoSerializer | 更高效的序列化 |
spark.sql.autoBroadcastJoinThreshold | 10-100MB | 自动广播 Join 阈值 |
spark.sql.files.maxPartitionBytes | 128MB | 单个分区文件大小上限 |
spark.sql.adaptive.coalescePartitions.enabled | true | 自动合并小分区 |
4.3 内存管理与 GC 调优
Executor 内存结构 (Unified Memory Management):
┌───────────────────────────────────────┐
│ Reserved Memory (300MB) │
├───────────────────────────────────────┤
│ User Memory (spark.memory.fraction) │
│ 用于存储用户数据结构、RDD transformations │
├───────────────────────────────────────┤
│ Spark Memory │
│ ┌─────────────┬─────────────────┐ │
│ │ Storage │ Execution │ │
│ │ (cache/persist) │ │
│ │ 默认 0.5 │ 默认 0.5 │ │
│ └─────────────┴─────────────────┘ │
└───────────────────────────────────────┘
# 高频 GC 场景调优
spark.conf.set("spark.memory.fraction", "0.8") # 给计算更多内存
spark.conf.set("spark.memory.storageFraction", "0.3") # 减少缓存占用
spark.conf.set("spark.executor.extraJavaOptions",
"-XX:+UseG1GC -XX:MaxGCPauseMillis=200")
5. 生产环境最佳实践
5.1 数据读写优化
# Parquet 优化
df.write \
.option("compression", "zstd") \
.option("parquet.block.size", "256MB") \
.mode("overwrite") \
.parquet("s3://output/")
# 分区策略
df.write.partitionBy("year", "month", "day").parquet("s3://output/")
# 分区字段 = 过滤字段,避免全表扫描
# 批量读取小文件
df = spark.read.option("mergeSchema", "true").parquet("s3://path/*")
# 或用与 Hive 配合的 ACID 表
5.2 Checkpoint 与容错
# Streaming 场景 Checkpoint(Structured Streaming)
query = streamDF.writeStream \
.format("parquet") \
.option("checkpointLocation", "s3://checkpoints/job1/") \
.start("s3://output/")
# RDD Checkpoint(批处理,截断 Lineage)
sparkContext.setCheckpointDir("hdfs:///checkpoints")
longRDD.checkpoint()
6. Spark 3.x 新特性详解
Spark 3.x 系列(3.0 ~ 3.5)引入了多项关键改进,使批处理性能、易用性和标准兼容性迈上新台阶。
6.1 Adaptive Query Execution (AQE)
AQE 是 Spark 3.0 引入的自适应查询执行框架,能在运行期根据真实统计信息动态调整执行计划,解决编译期统计信息不准导致的性能劣化问题。AQE 主要解决三大痛点:
- 自动合并 Shuffle 后的小分区:避免产生大量小文件和空任务
- 自动处理 Join 数据倾斜:将倾斜键拆分为多个子任务,均衡负载
- 动态切换 Join 策略:运行时检测表大小,将小表 Join 自动降级为 Broadcast Join
# AQE 完整生产配置
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.minPartitionSize", "1MB")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "400MB")
spark.conf.set("spark.sql.adaptive.localShuffleReader.enabled", "true")
spark.conf.set("spark.sql.adaptive.optimizeSkewsInRebalancePartitions.enabled", "true")
AQE 仅作用于 Exchange(Shuffle)节点之后,因此不会引入额外的 Stage 开销。建议在绝大多数生产环境默认开启。
6.2 动态分区裁剪 (Dynamic Partition Pruning, DPP)
DPP 解决了传统静态谓词下推无法跨越 Join 边界的限制。当大表(fact)与小表(dim)Join 时,Spark 会自动将小表的过滤结果广播到大表侧,用于动态剪枝分区文件:
-- DPP 自动生效,无需手动干预
SELECT f.*
FROM fact_orders f
JOIN dim_region d ON f.region_id = d.id
WHERE d.region = 'APAC';
-- 执行计划中会出现 "DynamicPruning" 标识
-- 等价于在 fact_orders 侧自动 injected filter: region_id IN (SELECT id FROM dim_region WHERE region = 'APAC')
DPP 生效条件:
- 被裁剪表必须是分区表
- Join 键与分区键一致
- 小表广播代价低于被裁剪的分区扫描代价
6.3 ANSI SQL 兼容模式
Spark 3.0 引入了 spark.sql.ansi.enabled,开启后行为与主流 ANSI SQL 更加一致:
- 数值溢出返回错误(而非静默回绕)
- 除零抛出
DIVIDE_BY_ZERO异常 CAST严格语义,非法格式报错而非返回NULL
# 开启 ANSI 模式,便于与下游数据仓库(Snowflake、Trino)保持一致语义
spark.conf.set("spark.sql.ansi.enabled", "true")
spark.conf.set("spark.sql.storeAssignmentPolicy", "ANSI") # 插入时严格类型检查
6.4 Pandas API on Spark
对于数据科学团队,Spark 3.2+ 提供的 pyspark.pandas(原 Koalas)实现了 90% 以上 Pandas API,将单机分析脚本无缝迁移到分布式环境:
import pyspark.pandas as ps
# 读取与 Pandas 完全一致
psdf = ps.read_parquet("s3://lake/sales/")
# 窗口函数、groupby、apply 自动分布式化
result = psdf.groupby("region").agg(
total_revenue=("amount", "sum"),
avg_order=("amount", "mean")
).sort_values("total_revenue", ascending=False)
# 与原生 Spark DataFrame 互转
spark_df = result.to_spark()
7. Spark SQL 深度优化
7.1 Catalyst 优化器原理
Catalyst 是 Spark SQL 的可扩展优化器,基于 Scala 的函数式编程特性,通过规则(Rule)和模式匹配(Pattern Matching)逐层变换执行计划。整个流程分为五个阶段:
- 解析(Analysis):将未解析的逻辑计划通过 Catalog 绑定元数据,解析表名、列名、类型
- 逻辑优化(Logical Optimization):应用 RBO(Rule-Based Optimization)规则,如谓词下推、列裁剪、常量折叠、连接重排序
- 物理规划(Physical Planning):使用 Cost Model 从候选物理计划中选择最优方案(CBO,基于表统计信息)
- 代码生成(Code Generation):Tungsten 将物理计划树编译为紧凑的 Java 字节码,消除虚函数调用
- 执行(Execution):生成 RDD DAG,提交到集群执行
// 使用 CBO(Cost-Based Optimization)前提:收集统计信息
spark.sql("ANALYZE TABLE orders COMPUTE STATISTICS FOR COLUMNS user_id, amount, dt")
spark.sql("ANALYZE TABLE orders COMPUTE STATISTICS NOSCAN")
// 查看完整优化过程
spark.sql("SELECT * FROM orders WHERE amount > 100 AND dt = '2024-01-01'")
.explain("cost") // 显示每个物理计划的估算代价
7.2 Tungsten 二进制格式与 UnsafeRow
Tungsten 绕开了 JVM 对象模型,直接在 off-heap 或堆内分配紧凑二进制数据(UnsafeRow)。其优势包括:
- 内存密度高:定长字段直接 inline,变长字段通过偏移量索引,消除对象头和对齐填充开销
- CPU 缓存友好:数据紧凑排布,顺序读取命中率高
- 零拷贝序列化:Shuffle 时无需反复序列化/反序列化 JVM 对象
7.3 Whole-Stage Codegen
传统火山模型(Volcano Iterator Model)中,每个算子通过 next() 链式调用,虚函数开销巨大。Whole-Stage Codegen 将一整个 Stage 内的所有算子融合为一个函数,使用 for 循环直接遍历数据:
// 查看 Codegen 生效的算子(Physical Plan 中显示 *)
spark.range(1000000)
.select($"id" * 2 + 1)
.filter($"id" > 100)
.groupBy($"id" % 10)
.agg(count("*"))
.explain()
// 输出示例:
// *(2) HashAggregate(keys=[(id#0L % 10)#3L], functions=[count(1)])
// +- *(2) HashAggregate ...
// +- *(1) Filter (id#0L > 100)
// +- *(1) Project [(id#0L * 2) + 1]
// +- *(1) Range (0, 1000000, step=1, splits=8)
// 其中 "*" 前缀表示该 Stage 已参与 Whole-Stage Codegen
当算子过于复杂(如包含非确定性表达式、Python UDF、外部数据源)时,Spark 会自动插入 WholeStageCodegen 边界,将可 codegen 部分与不可 codegen 部分隔离。
7.4 Broadcast Hash Join vs Sort-Merge Join 选择策略
| 特性 | Broadcast Hash Join | Sort-Merge Join |
|---|---|---|
| 适用条件 | 小表 <= autoBroadcastJoinThreshold(默认 10MB) | 两表均较大,无显著大小差异 |
| Shuffle 开销 | 无 Shuffle,小表广播到各节点 | 两表均 Shuffle,按 Join Key 排序 |
| 内存要求 | 小表需完整装入各 Executor 内存 | 内存压力较低,可 spill 到磁盘 |
| 倾斜容忍 | 不受倾斜影响(无 Shuffle) | 倾斜 Key 导致长尾任务 |
| 启动速度 | 快(无 Shuffle Stage) | 需要额外排序阶段 |
| AQE 介入 | 可选自动降级 | 支持 AQE 倾斜优化 |
// 强制使用 Sort-Merge Join(调优对比测试)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
// 强制广播 Join(显式 hint 或调大阈值)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "500MB")
// SQL HINT 方式
spark.sql("""
SELECT /*+ BROADCAST(dim) */ *
FROM fact JOIN dim ON fact.key = dim.key
""")
选型建议:
- 维表(< 100MB,实际大小经压缩后可广播)一律优先 Broadcast Hash Join
- 大表与大表关联,或无法预估小表大小时,使用 Sort-Merge Join + AQE 兜底
- 对于倾斜严重的大表 Join,AQE 的 Skew Join 优化优于手动加盐
8. Spark 内存管理
8.1 Unified Memory Model
Spark 1.6 之后采用统一内存管理模型(Unifed Memory Management),将 Execution Memory 与 Storage Memory 融合在同一区域,允许双方动态借用:
Executor JVM Heap
┌──────────────────────────────────────────────┐
│ Reserved Memory (300MB 固定) │
├──────────────────────────────────────────────┤
│ User Memory (1.0 - spark.memory.fraction) │
│ 用于用户自定义数据结构、Spark 内部元数据 │
├──────────────────────────────────────────────┤
│ Spark Memory (spark.memory.fraction, 默认0.6)│
│ ┌─────────────────┬─────────────────────┐ │
│ │ Storage Memory │ Execution Memory │ │
│ │ (persist/cache)│ (Shuffle/Sort/Join)│ │
│ │ 可溢出到磁盘 │ 可溢出到磁盘 │ │
│ │ 默认各占 0.5 │ 默认各占 0.5 │ │
│ │ 可被 Execution 借用(反之不可) │ │
│ └─────────────────┴─────────────────────┘ │
└──────────────────────────────────────────────┘
8.2 Storage Memory vs Execution Memory 仲裁机制
- Execution 优先于 Storage:当 Execution 内存不足时,可强制驱逐 Storage 内存中缓存的 RDD/DataFrame 分区
- Storage 不能抢占 Execution:Storage 空闲时 Execution 可借用,但 Execution 不会让出已占用的内存给 Storage
- 为什么要偏向 Execution? Shuffle 数据若无法内存计算,必须 spill 到磁盘,导致性能断崖式下跌;而缓存数据被驱逐后可以从数据源重读或根据 Lineage 重算
8.3 内存溢出诊断与调参
内存溢出(OOM)通常发生在以下场景:
- Executor 堆内存不足,大量对象无法 GC
spark.executor.memoryOverhead设置过低,Native 内存(Netty、PySpark、JNI)溢出导致容器被 K8s/YARN 强制 Kill(Exit Code 137)- 单个任务数据量过大(如
groupByKey后某个 Key 对应百万级记录)
# 生产环境内存调参模板
spark.conf.set("spark.executor.memory", "32g")
spark.conf.set("spark.executor.memoryOverhead", "8g") # 建议 memoryOverhead >= memory * 0.2
spark.conf.set("spark.executor.memoryFraction", "0.8") # Spark 3.x 中已废弃,仅适用于 2.x
spark.conf.set("spark.memory.fraction", "0.8") # 留给 Spark Memory 更多空间
spark.conf.set("spark.memory.storageFraction", "0.3") # 降低缓存比例,优先保障计算
# GC 调优(G1GC 为大堆推荐)
spark.conf.set(
"spark.executor.extraJavaOptions",
"-XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark"
)
# 诊断 OOM 时开启 verbose GC 日志(仅在调试期使用)
# -XX:+PrintGCDetails -XX:+PrintGCTimeStamps -Xloggc:/tmp/gc.log
诊断方法:
- Spark UI → Executors → 查看 “Task Time” 与 “GC Time” 比例,若 GC Time > 20% 说明堆内存吃紧
- YARN/ K8s 日志中出现
Container killed by YARN for exceeding memory limits时,增加memoryOverhead - 单个 Stage 内部部分任务执行时间远大于中位数(如 P99 是 P50 的 10 倍),通常是数据倾斜或内存溢出导致频繁 spill
9. Spark 数据倾斜治理
数据倾斜是生产环境最常见且最致命的瓶颈之一,表现为少数 Task 处理的数据量远大于其他 Task,导致 Stage 长尾。
9.1 倾斜识别与定位
通过 Spark UI 或 Spark History Server 进行诊断:
- Stage 页面:查看 “Event Timeline” 或 “Summary Metrics”,关注 “Duration” 的 Max 与 Median 差异
- Task 页面:按 “Shuffle Read Size” 或 “Input Size” 排序,若 Top N 任务数据量占总量的 80% 以上,则存在严重倾斜
- SQL / DAG 页面:查看 ShuffleExchange 节点后的聚合或 Join,定位具体算子
常见倾斜热点键:NULL 值、未登录用户的默认 user_id = 0、大客户的统一 merchant_id。
9.2 加盐打散(Salting)
加盐的核心思想是为倾斜键附加随机后缀,将单点热点拆散到多个分区进行局部聚合,随后再去盐完成全局聚合:
from pyspark.sql.functions import rand, concat, lit, col, sum as Fsum
# 假设 user_id = 0 是热点倾斜键
salt_count = 10 # 盐粒数
# 第一步:全局加盐聚合
df_with_salt = df.withColumn(
"salted_key",
concat(col("user_id"), lit("_"), (rand() * salt_count).cast("int"))
)
salted_agg = df_with_salt.groupBy("salted_key").agg(
Fsum("amount").alias("salted_sum")
)
# 第二步:去盐,恢复原始 Key 再全局聚合
from pyspark.sql.functions import split, element_at
final_agg = salted_agg.withColumn(
"user_id",
split(col("salted_key"), "_").getItem(0)
).groupBy("user_id").agg(
Fsum("salted_sum").alias("total_amount")
)
9.3 两阶段聚合
对于 RDD 或复杂 DataFrame 场景,两阶段聚合更为通用:
from pyspark.sql.functions import rand
# 阶段一:预聚合(带盐)
stage1 = df.withColumn("salt", (rand() * 10).cast("int")) \
.groupBy("user_id", "salt") \
.agg(sum("amount").alias("partial_sum"))
# 阶段二:全局聚合(无需再带盐)
stage2 = stage1.groupBy("user_id") \
.agg(sum("partial_sum").alias("total_amount"))
9.4 自定义 Partitioner 解决倾斜
当倾斜 Key 已知且分布固定时(例如 top 10 热门商品),可设计倾斜感知分区器,将热点均匀分散:
from pyspark import Partitioner
class SkewAwarePartitioner(Partitioner):
def __init__(self, num_partitions, skew_keys, replication=5):
self.num_partitions = num_partitions
self.skew_keys = set(skew_keys)
self.replication = replication
# 热点键映射到前 replication 个分区
def numPartitions(self):
return self.num_partitions
def getPartition(self, key):
if key in self.skew_keys:
return hash(key) % self.replication
return (hash(key) & 0x7fffffff) % (self.num_partitions - self.replication) + self.replication
# 使用
rdd = pair_rdd.partitionBy(SkewAwarePartitioner(200, skew_keys={"key1", "key2"}))
9.5 Skew Join 辅助手段对比
| 方法 | 适用场景 | 侵入性 | 性能影响 |
|---|---|---|---|
| AQE Skew Join | Spark 3.x、Shuffle Hash/Sort-Merge Join | 零侵入 | 轻量,自动 |
| 加盐打散 | 聚合操作、已知热点键 | 中等 | 增加一轮 Shuffle |
| 两阶段聚合 | RDD 复杂逻辑、非 SQL | 中等 | 增加一轮 Shuffle |
| 自定义 Partitioner | RDD、固定热点键集合 | 较高 | 需要在应用层维护键集合 |
| 广播 Join | 热点 Key 存在于小表侧 | 低 | 无 Shuffle,小内存代价 |
10. Spark 与 Delta Lake
Delta Lake 是构建在 Parquet 之上的开源存储层,为 Spark 提供了 ACID 事务、元数据管理和时间旅行能力,是数据湖向 Lakehouse 架构演进的核心组件。
10.1 ACID 事务保障
Delta Lake 通过事务日志(_delta_log)实现乐观并发控制:
from delta import DeltaTable
from pyspark.sql.functions import col, lit
# 条件更新(原子性)
delta_table = DeltaTable.forPath(spark, "s3://lake/orders")
delta_table.update(
condition=col("status") == "pending",
set={"status": lit("processed"), "updated_at": current_timestamp()}
)
# Merge(Upsert)操作
(delta_table.alias("target")
.merge(source_df.alias("source"), "target.order_id = source.order_id")
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.execute())
10.2 Time Travel 历史回溯
-- 查询历史版本(基于版本号)
SELECT * FROM delta.`s3://lake/orders` VERSION AS OF 15;
-- 查询历史版本(基于时间戳)
SELECT * FROM delta.`s3://lake/orders` TIMESTAMP AS OF '2024-06-01T00:00:00Z';
-- 查看版本历史
DESCRIBE HISTORY delta.`s3://lake/orders`;
10.3 Schema Enforcement & Evolution
Delta Lake 默认拒绝写入与表 Schema 不匹配的 DataFrame,避免脏数据污染数据湖:
# 严格 Schema 校验(默认行为)
df.write.format("delta").mode("append").save("s3://lake/orders")
# 自动 Schema Evolution(谨慎使用)
spark.conf.set("spark.databricks.delta.schema.autoMerge.enabled", "true")
df_new_columns.write.format("delta").mode("append").option("mergeSchema", "true").save("s3://lake/orders")
10.4 Z-Ordering 优化文件布局
Z-Ordering 通过多维度空间填充曲线(Z-order curve)对数据重新组织,使得多个常用过滤列的数据在物理文件上聚类,大幅减少 I/O:
-- 对 user_id 和 product_id 进行 Z-Order 优化(适合点查与范围过滤)
OPTIMIZE delta.`s3://lake/orders` ZORDER BY (user_id, product_id);
建议对高基数字段进行 Z-Ordering,低基数字段优先使用 Hive/Delta 分区。
10.5 Vacuum 过期文件清理
Delta Lake 的 Update/Delete/Merge 操作会产生历史版本文件,长期累积导致存储膨胀。Vacuum 用于清理不再被 Time Travel 引用的旧文件:
# 默认保留 7 天(168 小时),以下命令清理超过 7 天的旧版本文件
spark.conf.set("spark.databricks.delta.retentionDurationCheck.enabled", "false")
delta_table.vacuum(168) # 168 hours = 7 days
# 生产建议:配合 Airflow/Dagster 定时任务,每周运行一次 VACUUM
11. Spark 流批一体
Structured Streaming 将流处理抽象为在无限表(unbounded table)上的增量查询,复用与批处理完全一致的 DataFrame API。
11.1 微批 vs Continuous Processing 模式对比
| 特性 | Micro-Batch(默认) | Continuous Processing(实验性) |
|---|---|---|
| 延迟 | 毫秒级 ~ 秒级(默认 1s trigger) | 毫秒级 ~ 亚毫秒级 |
| 执行模型 | 周期性触发批作业 | 常驻长任务,持续处理 |
| 容错语义 | Exactly-once(Checkpoint + WAL) | Exactly-once(Checkpoint) |
| 适用算子 | 全部 SQL/DataFrame 算子 | 仅 Projection、Selection、Map、SQL Join(有限) |
| Source/Sink | 全面支持 | 仅 Kafka Source/Sink |
| 资源占用 | 每次 Trigger 调度开销 | 持续占用资源 |
| Spark 版本 | 稳定生产可用 | Spark 3.x 仍标注为实验特性 |
11.2 Watermark 与窗口聚合
Watermark 用于处理事件时间(Event Time)下的乱序数据,界定迟到数据的容忍窗口:
from pyspark.sql.functions import window, col, watermark, count, sum as Fsum
stream_df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "kafka:9092") \
.option("subscribe", "user_events") \
.option("startingOffsets", "latest") \
.load() \
.selectExpr("CAST(value AS STRING) as json") \
.select(from_json(col("json"), schema).alias("data")) \
.select("data.*")
# 定义 Watermark:允许事件迟到 10 分钟
windowed_counts = stream_df \
.withWatermark("event_time", "10 minutes") \
.groupBy(
window(col("event_time"), "5 minutes", "1 minute"), # 5分钟窗口,1分钟滑动步长
col("action")
) \
.agg(
count("*").alias("event_count"),
Fsum("value").alias("total_value")
)
Watermark 时间到达后,窗口状态才会被触发输出,并随后从 State Store 中清理以释放内存。
11.3 与 Kafka 集成
# Kafka Source
kafka_df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "broker1:9092,broker2:9092") \
.option("subscribe", "input-topic") \
.option("failOnDataLoss", "false") \
.option("maxOffsetsPerTrigger", 1000000) \
.load()
# Kafka Sink(至少需要一个 Key 或 Value 列)
query = windowed_counts \
.selectExpr("CAST(action AS STRING) as key", "to_json(struct(*)) AS value") \
.writeStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "broker1:9092,broker2:9092") \
.option("topic", "output-topic") \
.option("checkpointLocation", "s3://checkpoints/streaming-job/") \
.outputMode("update") \
.trigger(processingTime="10 seconds") \
.start()
query.awaitTermination()
12. Spark on Kubernetes
随着云原生架构普及,Spark on Kubernetes(K8s)逐渐取代 YARN 成为新一代资源调度方案。
12.1 spark-submit K8s 模式
# Cluster 模式提交(Driver 运行在 Pod 内)
spark-submit \
--master k8s://https://<k8s-api-server>:443 \
--deploy-mode cluster \
--name prod-batch-etl \
--class com.example.etl.BatchPipeline \
--conf spark.executor.instances=20 \
--conf spark.executor.cores=4 \
--conf spark.executor.memory=16g \
--conf spark.executor.memoryOverhead=4g \
--conf spark.kubernetes.container.image=registry/spark:3.5.0-scala2.12-java11 \
--conf spark.kubernetes.namespace=spark-jobs \
--conf spark.kubernetes.authenticate.driver.serviceAccountName=spark \
--conf spark.kubernetes.driver.pod.name=driver-prod-etl-$(date +%s) \
local:///opt/spark/jobs/etl-assembly.jar
12.2 Driver / Executor Pod 配置模板
# spark-pod-template.yaml
apiVersion: v1
kind: Pod
spec:
containers:
- name: spark-executor
resources:
requests:
memory: "16Gi"
cpu: "4"
limits:
memory: "20Gi"
cpu: "4"
env:
- name: AWS_REGION
value: "cn-north-1"
- name: S3_ACCESS_KEY
valueFrom:
secretKeyRef:
name: s3-credentials
key: access-key
volumeMounts:
- name: tmp-volume
mountPath: /tmp
volumes:
- name: tmp-volume
emptyDir:
medium: "Memory"
sizeLimit: "10Gi"
提交时通过 --conf spark.kubernetes.executor.podTemplateFile=/path/to/spark-pod-template.yaml 加载。
12.3 动态资源分配与调度器集成
# 动态资源分配(Dynamic Allocation)在 K8s 上需开启 shuffle tracking
spark.conf.set("spark.dynamicAllocation.enabled", "true")
spark.conf.set("spark.dynamicAllocation.shuffleTracking.enabled", "true")
spark.conf.set("spark.dynamicAllocation.minExecutors", "2")
spark.conf.set("spark.dynamicAllocation.maxExecutors", "100")
spark.conf.set("spark.dynamicAllocation.initialExecutors", "10")
spark.conf.set("spark.dynamicAllocation.executorIdleTimeout", "60s")
# Volcano / Yunikorn 调度器集成(支持队列与 Gang Scheduling)
spark.conf.set("spark.kubernetes.scheduler.name", "volcano")
spark.conf.set("spark.kubernetes.job.queue", "etl-queue")
13. 性能基准对比与选型矩阵
| 引擎 | 核心优势 | 最佳场景 | 劣势 |
|---|---|---|---|
| Spark | 批处理王者、生态完善、Lakehouse 原生 | 大规模 ETL、离线数仓、机器学习预处理 | 流处理延迟高于 Flink |
| Flink | 真正的流处理(毫秒级)、精确状态管理 | 实时监控、CEP、IoT 流处理 | 批处理生态弱于 Spark |
| Trino | 即席查询(Ad-hoc)、ANSI SQL 兼容好 | OLAP 交互式分析、联邦查询 | 无持久化状态,不适合复杂 ETL |
| 场景 | 推荐引擎 | 辅助技术 |
|---|---|---|
| 离线批处理 ETL(TB-PB 级) | Spark + Delta Lake | AQE、Z-Order、Kryo |
| 实时流处理(毫秒级延迟) | Flink | RocksDB State Backend |
| 近实时分析(分钟级延迟) | Spark Structured Streaming | Kafka + Delta Lake |
| Ad-hoc / BI 即席查询 | Trino / StarRocks | 物化视图、Connector 联邦查询 |
| 机器学习特征工程 | Spark MLlib + Pandas on Spark | Delta Lake Time Travel |
| 混合负载(Streaming + Batch) | Spark/Flink + 统一存储 | Delta Lake / Iceberg / Hudi |
14. 常见问题 (FAQ)
Q1: Spark SQL 中选择 Broadcast Join 时提示
BroadcastTimeout,应如何排查?
A: 首先确认小表实际大小是否超过spark.sql.autoBroadcastJoinThreshold,注意统计的广播大小是按列式压缩前的估算值。若小表确实较大,可尝试增加spark.sql.broadcastTimeout(默认 300s),或改用 Sort-Merge Join。若使用 AQE,可开启spark.sql.adaptive.enabled让 Spark 自动决策。
Q2: Executor 频繁被 K8s/YARN 以 Exit Code 137(OOM Killed)终止,如何定位真实原因?
A: Exit Code 137 通常不是 JVM 堆 OOM,而是进程整体内存(堆 + off-heap + Python/Py4J + Netty)超出容器限制。应调大spark.executor.memoryOverhead,比例建议不低于spark.executor.memory的 20%。同时检查是否有大广播变量、Python UDF 内存泄漏或 off-heap 缓存过大。
Q3: Spark UI 中某个 Stage 的 Max Task Duration 远高于 Avg/Median,是否一定是数据倾斜?
A: 绝大多数情况下是数据倾斜,但也可能是:单节点硬件故障(磁盘坏道、网卡降速)、G1GC 长时间 STW、或某个 Executor 上运行了其他抢占资源的进程。建议交叉对比 “Shuffle Read Size” / “Input Size” 指标,若数据量差异不大,则排查节点级问题。
Q4: Delta Lake 的 VACUUM 会不会误删还被 Time Travel 需要的文件?
A: VACUUM 默认只会删除超过 retention 周期(默认 7 天)的旧版本文件。如果你需要更长周期的审计回溯,应调大保留时间(如delta.logRetentionDuration = interval 30 days),并确保在 VACUUM 执行前已确认业务不再需要该时间段的历史版本。
Q5: Spark on Kubernetes 与 Spark on YARN 在生产环境如何抉择?
A: 若集群已深度使用 Hadoop 生态(HDFS、Hive on Tez)、且运维团队熟悉 YARN,优先保持 YARN。若正在向云原生迁移、需要更细粒度的资源隔离(命名空间、RBAC、Sidecar)、或与其他微服务共享 K8s 集群,则 Spark on K8s 是更优选择。短期共存策略:使用 K8s Operator(如 Spark Operator)统一管理 Spark 作业,降低两套资源调度器的运维复杂度。
总结
| 场景 | 推荐方案 |
|---|---|
| 通用批处理 | DataFrame + SQL |
| 需要编译期类型安全 | Dataset (Scala) |
| 非结构化/细粒度控制 | RDD |
| 聚合计算 | reduceByKey / aggregateByKey |
| 大小表 Join | Broadcast Hash Join(小表) / Sort-Merge Join(大表) |
| 数据倾斜 | AQE 自动优化 + 加盐打散 + 自定义 Partitioner |
| 生产调参 | Kryo + AQE + 合理分区数 + G1GC + 充足 memoryOverhead |
| 流批一体 | Structured Streaming + Watermark + Kafka 集成 |
| 云原生部署 | Spark on Kubernetes + 动态资源分配 + Volcano 调度 |
| 数据湖治理 | Delta Lake(ACID + Time Travel + Z-Order + Vacuum) |
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。