04. Apache Spark 批处理详解

深入 Apache Spark 批处理核心:RDD、DataFrame、Dataset 对比与选型,Spark SQL 查询优化,Shuffle 机制原理与调优策略,以及 Spark Core 关键参数配置。

Apache Spark 是目前最主流的分布式批处理引擎,其统一的编程模型和内存计算能力使其在大规模数据处理中占据核心地位。本文系统讲解 Spark 三大核心抽象、查询优化器、Shuffle 机制与生产调参策略。

1. Spark 核心抽象对比:RDD vs DataFrame vs Dataset

1.1 三大 API 全景对比

特性RDDDataFrameDataset
引入版本Spark 1.0Spark 1.3Spark 1.6
类型安全是(编译期)否(运行时检查)是(编译期)
模式推断
性能优化无内置优化Catalyst + TungstenCatalyst + Tungsten
API 风格函数式DSL/SQL函数式 + DSL
序列化Java 序列化Tungsten 二进制Tungsten 编码器
适用语言Scala/Java/Python/RScala/Java/Python/RScala/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 的算子groupByKeyreduceByKeyaggregateByKeysortByKeyjoincogrouprepartitiondistinct

3.2 Shuffle 优化策略

策略方法适用场景
减少 Shuffle 次数使用 reduceByKey 替代 groupByKey + map聚合操作
Map 端预聚合aggregateByKeycombineByKey键值对聚合
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.cores2-5每个 Executor 的 CPU 核数
spark.executor.memory8-32g堆内存大小
spark.executor.memoryOverheadmemory * 0.1-0.25off-heap/Native/Netty 内存
spark.sql.shuffle.partitions2-4倍 executor 总数Shuffle 后分区数
spark.default.parallelism2-3倍 executor 总数RDD 默认分区数
spark.sql.adaptive.enabledtrue自适应查询执行 (AQE)
spark.serializerKryoSerializer更高效的序列化
spark.sql.autoBroadcastJoinThreshold10-100MB自动广播 Join 阈值
spark.sql.files.maxPartitionBytes128MB单个分区文件大小上限
spark.sql.adaptive.coalescePartitions.enabledtrue自动合并小分区

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)逐层变换执行计划。整个流程分为五个阶段:

  1. 解析(Analysis):将未解析的逻辑计划通过 Catalog 绑定元数据,解析表名、列名、类型
  2. 逻辑优化(Logical Optimization):应用 RBO(Rule-Based Optimization)规则,如谓词下推、列裁剪、常量折叠、连接重排序
  3. 物理规划(Physical Planning):使用 Cost Model 从候选物理计划中选择最优方案(CBO,基于表统计信息)
  4. 代码生成(Code Generation):Tungsten 将物理计划树编译为紧凑的 Java 字节码,消除虚函数调用
  5. 执行(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 JoinSort-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 进行诊断:

  1. Stage 页面:查看 “Event Timeline” 或 “Summary Metrics”,关注 “Duration” 的 Max 与 Median 差异
  2. Task 页面:按 “Shuffle Read Size” 或 “Input Size” 排序,若 Top N 任务数据量占总量的 80% 以上,则存在严重倾斜
  3. 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 JoinSpark 3.x、Shuffle Hash/Sort-Merge Join零侵入轻量,自动
加盐打散聚合操作、已知热点键中等增加一轮 Shuffle
两阶段聚合RDD 复杂逻辑、非 SQL中等增加一轮 Shuffle
自定义 PartitionerRDD、固定热点键集合较高需要在应用层维护键集合
广播 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 LakeAQE、Z-Order、Kryo
实时流处理(毫秒级延迟)FlinkRocksDB State Backend
近实时分析(分钟级延迟)Spark Structured StreamingKafka + Delta Lake
Ad-hoc / BI 即席查询Trino / StarRocks物化视图、Connector 联邦查询
机器学习特征工程Spark MLlib + Pandas on SparkDelta 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
大小表 JoinBroadcast 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)

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 数据工程深度指南:Modern Data Stack 全栈实践
  2. 数据平台工程:Data Mesh、FinOps 与 DataOps 生产实践
  3. Kafka Connect CDC 实战:Debezium 数据同步与变更捕获