Spark 从零到进阶:数据开发工程师的系统学习指南

写在前面:如果你是一名刚接触 Spark 的数据开发工程师,面对网上零散的教程和概念感到无从下手,那么这篇文章就是为你准备的。我们将从"Spark 是什么"出发,一路走到性能调优与生产踩坑,配合大量代码示例和架构图,帮你建立完整的知识体系。本文基于 Spark 3.5.x 版本编写,并会提及 Spark 4.0 的预览方向。


目录

  1. Spark 概述
  2. Spark 架构与运行原理
  3. 环境搭建
  4. RDD 编程
  5. Spark SQL
  6. Spark Streaming vs Structured Streaming
  7. Spark 数据湖与 Lakehouse
  8. Spark 性能调优
  9. Spark 常见问题与踩坑
  10. 端到端实战项目
  11. Spark 3.x 新特性
  12. 学习路线与实战建议

1. Spark 概述

1.1 什么是 Spark

Apache Spark 是一个开源的统一分析引擎,专为大规模数据处理而设计。它最初于 2009 年在加州大学伯克利分校 AMPLab 诞生,2010 年开源,2014 年成为 Apache 顶级项目。Spark 提供了 SQL、流计算、机器学习和图计算等一整套大数据处理能力,可以运行在 Hadoop YARN、Kubernetes、Standalone 等多种集群管理器上。

1.2 发展历史

时间里程碑
2009Spark 在 UC Berkeley AMPLab 诞生
2010BSD 许可开源
2013捐赠给 Apache 软件基金会
2014Apache 顶级项目;Spark 1.0 发布
2016Spark 2.0:DataFrame/Dataset 成为主流 API,Structured Streaming
2020Spark 3.0:AQE、Dynamic Partition Pruning
2023Spark 3.4 / 3.5:Spark Connect 稳定、Pandas API 完善
2024+Spark 4.0 预览:更好的 ANSI SQL 兼容、性能持续优化

1.3 Spark vs Hadoop MapReduce

很多初学者会问"Spark 会替代 Hadoop 吗?"准确地说,Spark 替代的是 MapReduce 计算引擎,而 HDFS、YARN 依然是重要的存储和资源管理组件。

对比维度MapReduceSpark
中间结果落盘(HDFS)优先内存,可溢写磁盘
编程模型Map + Reduce 两阶段DAG(有向无环图),多阶段
延迟高(批处理)低(批流统一)
迭代计算每次迭代读写 HDFS内存缓存,迭代效率高 10-100x
生态仅 MapReduceSQL/Streaming/ML/Graph 一体化
容错基于磁盘复制RDD 血缘(Lineage)重算

1.4 核心特性

  • 速度快:DAG 执行引擎 + 内存计算,比 MapReduce 快 10-100 倍
  • 易用性:支持 Scala、Python、Java、R 四种语言 API,代码简洁
  • 通用性:一站式覆盖批处理、SQL、流计算、ML、图计算
  • 兼容性:可读写 HDFS、HBase、Cassandra、S3 等数据源,运行在 YARN/K8s/Standalone/Mesos 上

1.5 生态组件全景

在这里插入图片描述

  • Spark Core:底层引擎,提供 RDD 抽象、任务调度、内存管理、故障恢复
  • Spark SQL:结构化数据处理,DataFrame/Dataset API,兼容 Hive
  • Structured Streaming:基于 Spark SQL 的流计算引擎(原 Spark Streaming 已进入维护模式)
  • MLlib:分布式机器学习库
  • GraphX:图计算引擎

2. Spark 架构与运行原理

2.1 核心组件

在这里插入图片描述

组件职责
Driver运行 main() 函数,创建 SparkContext,将用户程序转化为 DAG,调度 Task
Executor在 Worker 节点上启动的 JVM 进程,执行 Task 并缓存数据
Cluster Manager集群资源管理器,负责分配 CPU/内存(YARN、K8s、Standalone)
Worker集群中可运行 Application 代码的节点

2.2 部署模式对比

模式特点适用场景
Local单机运行,多线程模拟分布式开发调试、学习
StandaloneSpark 自带集群管理器小规模独立集群
On YARN复用 Hadoop YARN 资源管理企业最常用,与 Hadoop 生态融合
On Kubernetes容器化部署,弹性扩缩容云原生环境,近年增长迅速

YARN 模式又分为 yarn-cluster(Driver 运行在 AM 中,适合生产)和 yarn-client(Driver 在客户端,适合交互调试)。

2.3 Job、Stage、Task 的划分

当一个 Action 算子被触发时,Spark 会提交一个 Job。DAGScheduler 根据 RDD 的依赖关系将 Job 划分为多个 Stage,每个 Stage 内创建一组 Task,由 TaskScheduler 分发到 Executor 执行。

在这里插入图片描述

划分规则:

  • 遇到宽依赖(Shuffle)就切分 Stage
  • 一个 Stage 内的所有 Task 执行完全相同的代码,只是处理的数据分区不同
  • Task 数量 = Stage 最后一个 RDD 的分区数

2.4 宽窄依赖与 Shuffle

类型说明典型算子
窄依赖父 RDD 每个分区最多被子 RDD 一个分区使用map、filter、union、mapPartitions
宽依赖父 RDD 每个分区被子 RDD 多个分区使用,需要 ShufflereduceByKey、groupByKey、join、distinct、repartition

在这里插入图片描述

Shuffle 过程:Map 端将数据按 Key 写入磁盘缓冲区,经过排序、合并、溢写后产生 Shuffle 文件;Reduce 端通过 BlockManager 拉取属于自己的数据。Shuffle 涉及磁盘 I/O、网络 I/O 和数据序列化,是 Spark 性能瓶颈的主要来源。

2.5 统一内存模型(Spark 1.6+)

Spark Executor 的内存分为以下区域:
在这里插入图片描述

  • Storage Memory:缓存 RDD、Broadcast 变量
  • Execution Memory:Shuffle、Join、Sort 等执行过程中的临时数据
  • 两者之间可以动态借用:Execution 可以借 Storage 的空闲内存,Storage 也可以借 Execution 的,但 Execution 优先级更高,Storage 被借走的部分在需要时会被淘汰

3. 环境搭建

3.1 Local 模式(最快上手)

# 下载(以 3.5.1 为例)
wget https://archive.apache.org/dist/spark/spark-3.5.1/spark-3.5.1-bin-hadoop3.tgz
tar -xzf spark-3.5.1-bin-hadoop3.tgz
cd spark-3.5.1-bin-hadoop3

# 启动 Spark Shell(Scala)
./bin/spark-shell --master local[*]

# 启动 PySpark
./bin/pyspark --master local[*]

local[*] 表示使用所有 CPU 核心。也可以用 local[2] 指定 2 个核心。

3.2 Standalone 模式

# 1. 配置 conf/spark-env.sh
cp conf/spark-env.sh.template conf/spark-env.sh
cat >> conf/spark-env.sh << 'EOF'
export SPARK_MASTER_HOST=master
export SPARK_WORKER_CORES=4
export SPARK_WORKER_MEMORY=8g
EOF

# 2. 配置 workers
echo "worker1" > conf/workers
echo "worker2" >> conf/workers

# 3. 启动集群
./sbin/start-all.sh

# 4. 提交应用
./bin/spark-submit \
  --master spark://master:7077 \
  --class com.example.MyApp \
  myapp.jar

3.3 On YARN 模式

确保 HADOOP_CONF_DIR 或 YARN_CONF_DIR 指向 Hadoop 配置目录:

export HADOOP_CONF_DIR=/etc/hadoop/conf

# 提交到 YARN(cluster 模式)
./bin/spark-submit \
  --master yarn \
  --deploy-mode cluster \
  --driver-memory 2g \
  --executor-memory 4g \
  --executor-cores 2 \
  --num-executors 10 \
  myapp.py

3.4 Spark Shell 与 spark-submit

工具用途
spark-shellScala 交互式 REPL,适合探索和调试
pysparkPython 交互式 REPL
spark-sqlSQL 交互式命令行
spark-submit提交应用到集群(生产环境)

spark-submit 常用参数:

参数说明示例
--master集群管理器 URLyarn / spark://host:7077 / local[*]
--deploy-modeDriver 运行位置cluster / client
--name应用名称MySparkApp
--classJava/Scala 主类com.example.MyApp
--jars额外依赖 JAR--jars lib/mysql-connector.jar
--packagesMaven 依赖自动下载--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1
--files分发文件到工作目录--files config.properties
--confSpark 配置项--conf spark.sql.shuffle.partitions=200
--driver-memoryDriver 内存4g
--executor-memory每个 Executor 内存8g
--executor-cores每个 Executor CPU 核数4
--num-executorsExecutor 数量(YARN)20
--total-executor-cores总核数(Standalone)100
--archives分发归档文件并解压--archives env.tar.gz#env
--py-files额外 Python 文件/压缩包--py-files utils.zip

3.5 PySpark 环境配置

# 方式一:pip 安装
pip install pyspark==3.5.1

# 方式二:使用 Spark 自带的 PySpark
export SPARK_HOME=/opt/spark
export PYTHONPATH=$SPARK_HOME/python:$SPARK_HOME/python/lib/py4j-0.10.9.7-src.zip:$PYTHONPATH
export PATH=$SPARK_HOME/bin:$PATH

# 验证
pyspark --version

在 Jupyter Notebook 中使用:

pip install jupyter
export PYSPARK_DRIVER_PYTHON=jupyter
export PYSPARK_DRIVER_PYTHON_OPTS='notebook'
pyspark --master local[*]

3.6 Spark on Kubernetes

Kubernetes(K8s)已成为大数据上云的主流部署方式。Spark on K8s 从 Spark 2.3 开始支持,Spark 3.1 进入 GA,目前在云原生场景下增长迅速。

部署模式原理:

Spark on K8s 将 Driver 和 Executor 直接作为 Pod 运行在 K8s 集群中:

在这里插入图片描述

  • Driver Pod 启动后,通过 K8s API 动态申请 Executor Pod
  • Executor Pod 运行结束后自动销毁,资源归还 K8s
  • 支持动态资源分配(Dynamic Allocation),按需扩缩容

与 YARN 模式对比:

对比维度Spark on YARNSpark on K8s
部署方式依赖 Hadoop YARN 集群容器化部署,镜像打包所有依赖
弹性扩缩容依赖 YARN 队列配置,扩容较慢秒级 Pod 创建,弹性能力强
资源隔离队列级隔离,Container 共享 OSPod 级隔离,容器边界清晰
依赖管理通过 --jars/--packages 分发打包进 Docker 镜像,版本一致性好
生态融合与 HDFS/Hive/HBase 深度集成云原生生态(Prometheus/Istio/Envoy)
运维复杂度Hadoop 运维体系成熟需 K8s 运维能力,学习曲线稍陡
多租户YARN 队列 + Linux 用户Namespace + RBAC,隔离更彻底
适用场景传统 Hadoop 集群、离线数仓云原生、混合云、弹性计算、流批一体

spark-submit on K8s 命令示例:

./bin/spark-submit \
  --master k8s://https://<k8s-apiserver-host>:<port> \
  --deploy-mode cluster \
  --name spark-pi \
  --class org.apache.spark.examples.SparkPi \
  --conf spark.executor.instances=5 \
  --conf spark.executor.memory=4g \
  --conf spark.executor.cores=2 \
  --conf spark.kubernetes.container.image=registry.example.com/spark:3.5.1 \
  --conf spark.kubernetes.namespace=spark-jobs \
  --conf spark.kubernetes.authenticate.driver.serviceAccountName=spark \
  --conf spark.kubernetes.driver.pod.name=spark-pi-driver \
  local:///opt/spark/examples/jars/spark-examples_2.12-3.5.1.jar

Docker 镜像构建要点:

# 基础镜像
FROM openjdk:11-jre-slim

# 安装 Spark
ARG SPARK_VERSION=3.5.1
ARG HADOOP_VERSION=3
RUN apt-get update && apt-get install -y curl tini && \
    curl -sL https://archive.apache.org/dist/spark/spark-${SPARK_VERSION}/spark-${SPARK_VERSION}-bin-hadoop${HADOOP_VERSION}.tgz | tar xz -C /opt && \
    ln -s /opt/spark-${SPARK_VERSION}-bin-hadoop${HADOOP_VERSION} /opt/spark && \
    apt-get clean

# 安装 Python(PySpark 场景)
RUN apt-get install -y python3 python3-pip && \
    pip3 install pyspark==${SPARK_VERSION} pandas pyarrow

# 拷贝自定义 JAR 依赖(如 JDBC 驱动、数据湖 SDK)
COPY jars/mysql-connector-j-8.0.33.jar /opt/spark/jars/
COPY jars/iceberg-spark-runtime-3.5_2.12-1.5.2.jar /opt/spark/jars/

# 拷贝应用 JAR
COPY app/your-app.jar /opt/spark/user-jars/

ENTRYPOINT ["/usr/bin/tini", "--"]

构建并推送镜像:

docker build -t registry.example.com/spark:3.5.1-custom .
docker push registry.example.com/spark:3.5.1-custom

K8s 特有配置参数:

参数说明示例
spark.kubernetes.container.imageDocker 镜像地址registry/spark:3.5.1
spark.kubernetes.namespaceK8s 命名空间spark-jobs
spark.kubernetes.authenticate.driver.serviceAccountNameDriver 使用的 ServiceAccountspark
spark.kubernetes.driver.pod.nameDriver Pod 名称spark-app-driver
spark.kubernetes.executor.podNamePrefixExecutor Pod 名称前缀spark-app-exec
spark.kubernetes.driver.limit.coresDriver CPU 限制2
spark.kubernetes.executor.limit.coresExecutor CPU 限制4
spark.kubernetes.driver.request.coresDriver CPU 请求1
spark.kubernetes.executor.request.coresExecutor CPU 请求2
spark.kubernetes.memoryOverheadFactor堆外内存比例0.2
spark.kubernetes.allocation.batch.size每批申请 Pod 数10
spark.kubernetes.executor.deleteOnTermination完成后删除 Executor Podtrue
spark.kubernetes.driver.podTemplateFileDriver Pod 模板文件driver-pod.yaml
spark.kubernetes.executor.podTemplateFileExecutor Pod 模板文件executor-pod.yaml
spark.kubernetes.file.upload.path应用依赖上传路径s3://spark-uploads/
spark.kubernetes.authenticate.caCertFileCA 证书路径/var/run/secrets/...

适用场景:

  • 云原生架构:与公司现有 K8s 平台统一运维,避免 Hadoop/YARN 独立集群
  • 弹性扩缩容:利用 K8s Cluster Autoscaler,根据负载自动增减节点,降低成本
  • 混合云/多云:同一镜像在不同云厂商 K8s 集群运行,避免厂商锁定
  • 流批一体:Streaming 长任务与离线批任务共享 K8s 资源池,通过 Namespace/Quota 隔离
  • 依赖隔离:不同应用使用不同镜像,彻底解决"一个集群多版本依赖冲突"问题

💡 生产建议:K8s 上建议启用 External Shuffle Service(或使用 Spark 3.0+ 的 Shuffle Data on PVC 方案),开启 Dynamic Allocation,并配置 Pod 模板以挂载共享存储(如 HDFS/S3/OSS)。


4. RDD 编程

虽然现在 Spark SQL/DataFrame 是主流 API,但理解 RDD(Resilient Distributed Dataset,弹性分布式数据集)是掌握 Spark 底层原理的关键。

4.1 RDD 五大特性

  1. 分区列表(Partitions):数据被切分为多个分区,每个分区在一个节点上计算
  2. 计算函数(Compute):每个分区都有一个计算函数来生成数据
  3. 依赖关系(Dependencies):RDD 之间有血缘关系,用于故障恢复
  4. 分区器(Partitioner):KV 类型 RDD 可选(Hash/Range),决定数据分布
  5. 优先位置(Preferred Locations):每个分区的计算优先调度到数据所在节点

4.2 创建 RDD

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("RDDDemo").master("local[*]").getOrCreate()
sc = spark.sparkContext

# 方式一:从集合创建
rdd = sc.parallelize([1, 2, 3, 4, 5], numSlices=3)

# 方式二:从外部存储创建
rdd = sc.textFile("hdfs:///data/input.txt")
rdd = sc.wholeTextFiles("hdfs:///data/logs/")  # 读取目录下所有文件

4.3 Transformation vs Action

Transformation 是懒执行的,只记录转换逻辑,不立即计算;Action 触发真正的计算。

类型常见算子
Transformationmap、filter、flatMap、mapPartitions、sample、union、intersection、distinct、groupByKey、reduceByKey、aggregateByKey、sortByKey、join、cogroup、cartesian、coalesce、repartition
Actioncollect、count、take、first、takeOrdered、reduce、fold、aggregate、foreach、saveAsTextFile、countByKey、collectAsMap
# Transformation 示例
nums = sc.parallelize([1, 2, 3, 4, 5, 6])

# map:一对一转换
squares = nums.map(lambda x: x * x)        # [1, 4, 9, 16, 25, 36]

# filter:过滤
evens = nums.filter(lambda x: x % 2 == 0)  # [2, 4, 6]

# flatMap:一对多展平
lines = sc.parallelize(["hello world", "hello spark"])
words = lines.flatMap(lambda line: line.split(" "))
# ["hello", "world", "hello", "spark"]

# reduceByKey:按 Key 聚合(Map 端预聚合,性能优于 groupByKey)
pairs = words.map(lambda w: (w, 1))
counts = pairs.reduceByKey(lambda a, b: a + b)
# [("hello", 2), ("world", 1), ("spark", 1)]

# join
rdd1 = sc.parallelize([("a", 1), ("b", 2)])
rdd2 = sc.parallelize([("a", "x"), ("b", "y"), ("b", "z")])
rdd1.join(rdd2).collect()
# [("a", (1, "x")), ("b", (2, "y")), ("b", (2, "z"))]

# sortByKey
sorted_counts = counts.sortByKey(ascending=True)

# Action 示例
print(counts.collect())      # 返回所有结果到 Driver
print(counts.count())        # 元素数量
print(counts.take(2))        # 取前 2 个
print(counts.first())        # 第一个
total = nums.reduce(lambda a, b: a + b)  # 聚合
counts.saveAsTextFile("hdfs:///output/wordcount")

4.4 reduceByKey vs groupByKey

这是面试高频考点:

对比reduceByKeygroupByKey
Map 端预聚合✅ 有(Combiner)❌ 无
Shuffle 数据量小大
性能好差
适用场景聚合类操作需要遍历所有值的场景

优先使用 reduceByKey / aggregateByKey / foldByKey,它们会在 Map 端做预聚合。

4.5 持久化:cache / persist / checkpoint

from pyspark import StorageLevel

# cache() = persist(MEMORY_ONLY)
rdd.cache()

# persist 支持多种存储级别
rdd.persist(StorageLevel.MEMORY_AND_DISK)      # 内存不够溢写磁盘
rdd.persist(StorageLevel.MEMORY_ONLY_SER)      # 序列化后存内存,节省空间
rdd.persist(StorageLevel.MEMORY_AND_DISK_SER_2) # 序列化 + 磁盘 + 2 副本

# 解除缓存
rdd.unpersist()

# checkpoint:将 RDD 写入 HDFS 等可靠存储,切断血缘
sc.setCheckpointDir("hdfs:///checkpoint")
rdd.checkpoint()
方式存储位置血缘适用场景
cache内存保留频繁重用的小数据
persist内存/磁盘可选保留根据数据量选择
checkpointHDFS切断长血缘链、迭代计算

4.6 共享变量

默认情况下,Spark 会把函数中用到的变量拷贝到每个 Task,Task 对变量的修改不会回传 Driver。共享变量提供了两种特殊机制:

Broadcast Variable(广播变量):将只读变量缓存到每个 Executor,而不是每个 Task 一份,大幅减少数据传输。

# 大数据量的查找表
lookup_table = {"US": "United States", "CN": "China", "JP": "Japan"}
bc_var = sc.broadcast(lookup_table)

rdd = sc.parallelize(["US", "CN", "JP"])
result = rdd.map(lambda code: bc_var.value.get(code, code)).collect()
# ["United States", "China", "Japan"]

Accumulator(累加器):分布式计数器,只能在 Executor 端 add,在 Driver 端读取 value。

accum = sc.accumulator(0)

def process_line(line):
    global accum
    if "ERROR" in line:
        accum.add(1)
    return line

sc.textFile("hdfs:///logs/app.log").map(process_line).count()
print(f"Error count: {accum.value}")

⚠️ 注意:Accumulator 在 Transformation 中可能因 Task 重试导致重复计数,建议在 Action(如 foreach)中使用,或使用 Named Accumulator 并在 Driver 端确认只读取一次。


5. Spark SQL

Spark SQL 是当前 Spark 最核心、使用最广泛的模块。DataFrame/Dataset API 比 RDD 更高级、更优化。

5.1 DataFrame 与 Dataset

特性RDDDataFrameDataset
数据模型无结构带 Schema 的行带 Schema 的强类型对象
类型安全是(编译时)否(运行时)是(编译时)
优化无CatalystCatalyst
语言Scala/Java/Python/R全部主要 Scala/Java
Tungsten否是是

在 PySpark 中,DataFrame 是主要 API(Python 没有 Dataset 的编译时类型检查)。

5.2 创建 DataFrame

from pyspark.sql import SparkSession
from pyspark.sql.types import *
from pyspark.sql.functions import col, sum as _sum, avg, count

spark = SparkSession.builder \
    .appName("SparkSQLDemo") \
    .master("local[*]") \
    .config("spark.sql.shuffle.partitions", "4") \
    .getOrCreate()

# 1. 从 RDD 转换
rdd = sc.parallelize([(1, "Alice", 25), (2, "Bob", 30)])
df = spark.createDataFrame(rdd, ["id", "name", "age"])

# 2. 显式指定 Schema
schema = StructType([
    StructField("id", IntegerType(), False),
    StructField("name", StringType(), True),
    StructField("age", IntegerType(), True),
])
df = spark.createDataFrame(rdd, schema)

# 3. 读取文件
df = spark.read.csv("data/users.csv", header=True, inferSchema=True)
df = spark.read.json("data/users.json")
df = spark.read.parquet("data/users.parquet")
df = spark.read.orc("data/users.orc")

# 4. 读取 Hive 表
df = spark.sql("SELECT * FROM default.users")

# 5. JDBC
df = spark.read.format("jdbc") \
    .option("url", "jdbc:mysql://localhost:3306/mydb") \
    .option("dbtable", "users") \
    .option("user", "root") \
    .option("password", "xxx") \
    .load()

# 6. toDF
df = [(1, "Alice"), (2, "Bob")].toDF(["id", "name"])  # Scala 风格
# PySpark 中:
df = spark.createDataFrame([(1, "Alice"), (2, "Bob")], ["id", "name"])

5.3 常用 Transformation

# 选择列
df.select("name", "age").show()
df.select(col("name"), col("age") + 1).show()

# 过滤
df.filter(col("age") > 25).show()
df.where("age > 25").show()

# 分组聚合
df.groupBy("department").agg(
    count("*").alias("emp_count"),
    _sum("salary").alias("total_salary"),
    avg("salary").alias("avg_salary")
).show()

# 排序
df.orderBy(col("age").desc()).show()

# 去重
df.select("department").distinct().show()

# 列重命名
df.withColumnRenamed("name", "employee_name")

# 新增列
df.withColumn("age_next_year", col("age") + 1)

# Join
df1.join(df2, on="id", how="inner")           # inner/left/right/outer/semi/anti
df1.join(df2, df1.id == df2.emp_id, "left_outer")

# 窗口函数
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, rank, dense_rank

window_spec = Window.partitionBy("department").orderBy(col("salary").desc())
df.withColumn("rank", row_number().over(window_spec)).show()

# 写数据
df.write.mode("overwrite").parquet("hdfs:///output/users")
df.write.partitionBy("department").parquet("hdfs:///output/users_by_dept")

5.4 UDF / UDAF

from pyspark.sql.functions import udf
from pyspark.sql.types import StringType, IntegerType

# 标量 UDF
def age_group(age):
    if age < 25:
        return "Young"
    elif age < 40:
        return "Mid"
    else:
        return "Senior"

age_group_udf = udf(age_group, StringType())
df.withColumn("age_group", age_group_udf(col("age"))).show()

# 更高效的方式:pandas UDF(向量化执行,基于 Apache Arrow)
import pandas as pd
from pyspark.sql.functions import pandas_udf

@pandas_udf(StringType())
def age_group_pandas(age: pd.Series) -> pd.Series:
    return age.apply(age_group)

df.withColumn("age_group", age_group_pandas(col("age"))).show()

💡 性能提示:普通 Python UDF 每行一次 Python/JVM 序列化,性能差;Pandas UDF 以列式批量传输,性能提升 10-100 倍。Spark 3.5 还引入了 Python UDF profiling 和更高效的 Arrow 批处理。

5.5 开窗函数

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, rank, dense_rank, lag, lead

window_spec = Window.partitionBy("department").orderBy(col("salary").desc())

# 排名
df.select(
    "name", "department", "salary",
    row_number().over(window_spec).alias("rn"),
    rank().over(window_spec).alias("rank"),
    dense_rank().over(window_spec).alias("dr"),
    lag("salary", 1).over(window_spec).alias("prev_salary"),
    lead("salary", 1).over(window_spec).alias("next_salary"),
).show()

5.6 执行计划与 Catalyst 优化器

df.explain()          # 简明物理计划
df.explain(True)      # 含解析/分析/优化/物理计划
df.explain("formatted")  # 格式化输出(Spark 3.0+)

Catalyst 优化器是 Spark SQL 的核心,它将用户的 DataFrame/SQL 代码经过以下阶段优化:

在这里插入图片描述

常见优化规则:

  • 谓词下推(Predicate Pushdown):将过滤条件尽早下推到数据源
  • 列裁剪(Column Pruning):只读取需要的列
  • 常量折叠(Constant Folding):预计算常量表达式
  • 布尔简化:简化 AND/OR 逻辑

5.7 Tungsten 优化

Tungsten 是 Spark 的执行引擎优化项目,目标是最大化 CPU 效率和内存效率:

  1. 内存管理:使用堆外内存(off-heap),避免 JVM GC 开销
  2. 二进制处理:数据以二进制格式存储,避免 Java 对象的序列化/反序列化
  3. Whole-Stage CodeGen:将整个 Stage 的多个算子融合为一个 Java 函数,消除虚函数调用
  4. 向量化读取:Parquet/ORC 列式批量读取

5.8 AQE(Adaptive Query Execution,自适应查询执行)

Spark 3.0 引入的重大特性,在 Shuffle Map 阶段完成后根据运行时统计信息动态调整执行计划:

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 分区自动将过小的分区合并,减少 Task 数量
动态处理数据倾斜自动检测倾斜分区并拆分 Join
动态切换 Join 策略运行时发现小表可广播时,自动从 Sort-Merge Join 切换为 Broadcast Join

6. Spark Streaming vs Structured Streaming

6.1 Spark Streaming(DStream)

Spark Streaming 是早期的流计算方案,基于微批处理(Micro-Batch),将实时数据流按时间间隔切分为小的 RDD 批处理。
在这里插入图片描述

from pyspark.streaming import StreamingContext

ssc = StreamingContext(sc, batchDuration=5)  # 5 秒一个批次
lines = ssc.socketTextStream("localhost", 9999)
counts = lines.flatMap(lambda line: line.split(" ")) \
              .map(lambda w: (w, 1)) \
              .reduceByKey(lambda a, b: a + b)
counts.pprint()
ssc.start()
ssc.awaitTermination()

⚠️ Spark Streaming(DStream)从 Spark 3.4 起已标记为维护模式,新项目应使用 Structured Streaming。

6.2 Structured Streaming

Structured Streaming 是基于 Spark SQL 引擎构建的端到端流计算,将实时数据流视为一张持续追加的"无界表"。

from pyspark.sql.functions import *

# 从 Kafka 读取
df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "orders") \
    .load()

# 解析 JSON
parsed = df.selectExpr("CAST(value AS STRING) as json") \
    .select(from_json("json", schema).alias("data")) \
    .select("data.*")

# 聚合
agg = parsed.groupBy("product_id").agg(
    count("*").alias("order_count"),
    sum("amount").alias("total_amount")
)

# 输出到控制台
query = agg.writeStream \
    .outputMode("complete") \
    .format("console") \
    .trigger(processingTime="10 seconds") \
    .start()

query.awaitTermination()

6.3 核心概念

Event Time(事件时间):数据产生时自带的时间戳,而非 Spark 处理时间(Processing Time)。

Watermark(水位线):定义了系统等待迟到数据的时间阈值,超过阈值的迟到数据将被丢弃。

# 设置 10 分钟 watermark
df_with_watermark = parsed \
    .withWatermark("event_time", "10 minutes") \
    .groupBy(
        window("event_time", "5 minutes"),  # 5 分钟滚动窗口
        "product_id"
    ).agg(count("*").alias("cnt"))

Window(窗口):

窗口类型说明
Tumbling Window固定大小,不重叠(如每 5 分钟统计)
Sliding Window固定大小,可重叠(如每 1 分钟统计最近 5 分钟)
Session Window基于活动间隔动态划分(Spark 3.2+)

Output Mode(输出模式):

模式说明适用聚合
Append仅输出新行无聚合
Complete输出全量结果有聚合
Update仅输出更新行有聚合(Spark 2.1+)

6.4 Source 与 Sink

Source说明
Kafka最常用,支持从 Kafka 读取消息
File监听目录中新文件
Socket测试用,从 TCP Socket 读取
Rate测试用,每秒生成指定行数
Sink说明
Kafka写入 Kafka Topic
File写入文件(Parquet/JSON/CSV)
Console控制台(调试用)
Foreach/ForeachBatch自定义写入逻辑
Memory存储为内存表(调试用)

6.5 Exactly-Once 语义

Structured Streaming 通过以下机制保证精确一次(Exactly-Once):

  1. 可重放的 Source:如 Kafka,记录 offset 可重新读取
  2. 幂等的 Sink:或使用事务写入(如 Kafka 事务、文件原子写入)
  3. Checkpoint + WAL:将 offset 和聚合状态持久化到可靠存储
query = agg.writeStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("topic", "order_summary") \
    .option("checkpointLocation", "hdfs:///checkpoint/orders") \
    .outputMode("complete") \
    .start()

注意:foreachBatch 本身不保证 Exactly-Once,需要 Sink 端实现幂等或使用事务。

7. Spark 数据湖与 Lakehouse

随着大数据从"数据仓库"向"湖仓一体"演进,数据湖格式已成为 Spark 生态不可或缺的一环。本章介绍三大开源数据湖格式及 Paimon,并给出 Spark 实战代码。

7.1 数据湖概念与三剑客对比

数据湖(Data Lake):以原始格式存储海量结构化、半结构化、非结构化数据的集中式存储,支持多种计算引擎直接读写。传统数据湖缺乏事务支持,容易产生脏数据。

Lakehouse(湖仓一体):在数据湖存储之上增加事务层、Schema 管理、ACID 语义,兼具数据湖的灵活性和数据仓库的管理能力。

特性Delta LakeApache IcebergApache Hudi
定位Databricks 主导,强绑定 Spark通用表格式,引擎无关流式数据湖,Upsert 优先
ACID 事务✅✅✅
Schema 演化✅(新增/删除列)✅(新增/删除/重命名/调序)✅
分区演化❌✅(隐式分区演化)✅
时间旅行✅✅(Snapshot ID/时间戳)✅
Upsert/Merge✅ MERGE INTO✅ MERGE INTO✅(原生 Upsert 是核心特性)
删除更新✅✅(V2 Row-Level)✅(COW/MOR)
流写入✅✅(Spark Structured Streaming)✅(深度流式支持)
计算引擎Spark 为主(Delta 2.0+ 支持 Flink/Presto)Spark/Flink/Trino/Presto/HiveSpark/Flink/Trino/Presto
社区活跃度Databricks 商业驱动Apache 顶级项目,中立开放Apache 顶级项目,Uber 起源
典型用户Databricks 客户Apple/Netflix/LinkedInUber/Grab/字节跳动
特色Z-Order 优化、Delta UniForm隐藏分区、分区演化、快照隔离COW/MOR 双模式、Compaction/Clustering

Apache Paimon 简介(原 Flink Table Store):

Paimon 是 Flink 社区发起的流批一体存储项目,2024 年成为 Apache 顶级项目。它的核心定位是流式数据湖:

  • 原生支持 Flink CDC 摄入,毫秒级延迟 Upsert
  • 同时支持 Spark 读写,流批统一
  • LSM Tree 架构,高吞吐写入与查询兼顾
  • 适合实时数仓场景(Flink + Paimon + Spark 分析)

选型提示:如果团队以 Spark 为核心且用 Databricks,Delta Lake 最省心;如果追求引擎中立和多引擎生态,Iceberg 是趋势;如果核心需求是流式 Upsert 和增量处理,Hudi/Paimon 更合适。

7.2 Spark + Iceberg 实战

Iceberg 核心概念:

概念说明
Snapshot(快照)表的一次完整状态,每次写入生成新快照,旧快照保留用于时间旅行
Manifest(清单)记录数据文件路径、分区信息、列级统计,Manifest List 管理多个 Manifest
Metadata File表元数据入口,记录当前 Snapshot、Schema、Partition Spec 等
Partition Spec分区规则,支持隐藏分区(如按天分区但查询可按小时过滤)
Schema Evolution新增/删除/重命名/调整列顺序,无需重写数据文件
Partition Evolution修改分区规则后旧数据不受影响,新数据按新规则写入

Spark 提交 Iceberg 依赖配置:

spark-submit \
  --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.5.2 \
  --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
  --conf spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkSessionCatalog \
  --conf spark.sql.catalog.spark_catalog.type=hive \
  --conf spark.sql.catalog.local=org.apache.iceberg.spark.SparkCatalog \
  --conf spark.sql.catalog.local.type=hadoop \
  --conf spark.sql.catalog.local.warehouse=s3://my-bucket/iceberg-warehouse \
  your_app.py

建表与写入:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, current_date, date_format

spark = SparkSession.builder.appName("IcebergDemo").getOrCreate()

# 建表
spark.sql("""
CREATE TABLE IF NOT EXISTS local.db.orders (
    order_id    BIGINT,
    user_id     BIGINT,
    product_id  BIGINT,
    amount      DECIMAL(10,2),
    order_time  TIMESTAMP,
    dt          STRING
) USING iceberg
PARTITIONED BY (dt)
""")

# 批量写入
batch_df = spark.read.parquet("s3://raw-data/orders/2024-01-15/")
batch_df.writeTo("local.db.orders").overwritePartitions()

# 流式写入(Structured Streaming)
from pyspark.sql.functions import from_json, col

kafka_schema = "order_id BIGINT, user_id BIGINT, product_id BIGINT, " \
               "amount DECIMAL(10,2), order_time TIMESTAMP, dt STRING"

stream_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "orders") \
    .load() \
    .selectExpr("CAST(value AS STRING) as json") \
    .select(from_json(col("json"), kafka_schema).alias("data")) \
    .select("data.*")

stream_df.writeStream \
    .format("iceberg") \
    .option("path", "local.db.orders") \
    .option("checkpointLocation", "s3://checkpoints/orders") \
    .trigger(processingTime="1 minute") \
    .start()

查询与时间旅行:

# 查询当前快照
spark.sql("SELECT * FROM local.db.orders WHERE dt = '2024-01-15'").show()

# Time Travel:按 Snapshot ID
spark.sql("""
SELECT * FROM local.db.orders.snapshotId(12345678901234567)
""").show()

# Time Travel:按时间戳
spark.sql("""
SELECT * FROM local.db.orders.timestampAsOf('2024-01-15 10:00:00')
""").show()

# 查看快照历史
spark.sql("SELECT * FROM local.db.orders.snapshots").show()

# 回滚到指定快照
spark.sql("CALL local.system.rollback_to_snapshot('db.orders', 12345678901234567)")

Schema 演化:

# 新增列
spark.sql("ALTER TABLE local.db.orders ADD COLUMN status STRING")

# 删除列
spark.sql("ALTER TABLE local.db.orders DROP COLUMN status")

# 重命名列
spark.sql("ALTER TABLE local.db.orders RENAME COLUMN amount TO total_amount")

# 调整列顺序
spark.sql("ALTER TABLE local.db.orders ALTER COLUMN order_id AFTER order_time")

分区演化:

# 从按天分区改为按月分区(旧数据不受影响)
spark.sql("ALTER TABLE local.db.orders ADD PARTITION FIELD months(order_time)")
spark.sql("ALTER TABLE local.db.orders DROP PARTITION FIELD dt")

MERGE INTO(Upsert):

# 创建更新数据集
updates_df = spark.createDataFrame([
    (1001, 2001, 3001, 99.99, "2024-01-15 14:30:00", "2024-01-15", "PAID"),
    (1002, 2002, 3002, 59.99, "2024-01-15 15:00:00", "2024-01-15", "PAID"),
], ["order_id", "user_id", "product_id", "amount", "order_time", "dt", "status"])

updates_df.createOrReplaceTempView("updates")

spark.sql("""
MERGE INTO local.db.orders t
USING updates s
ON t.order_id = s.order_id
WHEN MATCHED THEN
    UPDATE SET t.amount = s.amount, t.status = s.status, t.order_time = s.order_time
WHEN NOT MATCHED THEN
    INSERT (order_id, user_id, product_id, amount, order_time, dt, status)
    VALUES (s.order_id, s.user_id, s.product_id, s.amount, s.order_time, s.dt, s.status)
""")

7.3 Spark + Hudi 实战

Hudi 表类型对比:

特性COW(Copy On Write)MOR(Merge On Read)
写入方式每次写入更新重写整个 Parquet 文件增量写入 Log 文件,定期合并
写入延迟高(重写文件)低(追加日志)
读取延迟低(直接读 Parquet)稍高(需合并 Parquet + Log)
Compaction不需要需要(同步/异步)
适用场景批量更新、读多写少流式 Upsert、写多读多
存储效率高(列式压缩)稍低(Log 行式存储)
典型场景离线 ETL、T+1 报表实时数仓、近实时分析

Hudi 核心概念:

概念说明
Timeline表的操作历史,包含 Commit/DeltaCommit/Compaction/Clean,每个 Instant 记录状态
File Group一组文件,由 File ID 标识,包含 Base File(Parquet)和 Log File
File Slice某一时刻 File Group 的快照 = 一个 Base File + 对应 Log 文件
CompactionMOR 表将 Log 文件合并到 Base File 的过程,可同步或异步执行
Clustering将小文件合并、按列聚簇重写,优化查询性能

Spark 写入 Hudi 代码示例:

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("HudiDemo") \
    .config("spark.jars.packages",
            "org.apache.hudi:hudi-spark3.5-bundle_2.12:0.15.0") \
    .config("spark.sql.extensions",
            "org.apache.spark.sql.hudi.HoodieSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog",
            "org.apache.spark.sql.hudi.catalog.HoodieCatalog") \
    .getOrCreate()

table_path = "s3://my-bucket/hudi/orders"
table_name = "orders"

# 批量导入(bulk_insert,最高效,不做去重)
df = spark.read.parquet("s3://raw-data/orders/")
df.write.format("hudi") \
    .option("hoodie.table.name", table_name) \
    .option("hoodie.datasource.write.recordkey.field", "order_id") \
    .option("hoodie.datasource.write.partitionpath.field", "dt") \
    .option("hoodie.datasource.write.precombine.field", "order_time") \
    .option("hoodie.datasource.write.operation", "bulk_insert") \
    .mode("overwrite") \
    .save(table_path)

# Upsert(默认操作,按主键更新或插入)
updates_df = spark.read.parquet("s3://raw-data/orders/incremental/")
updates_df.write.format("hudi") \
    .option("hoodie.table.name", table_name) \
    .option("hoodie.datasource.write.recordkey.field", "order_id") \
    .option("hoodie.datasource.write.partitionpath.field", "dt") \
    .option("hoodie.datasource.write.precombine.field", "order_time") \
    .option("hoodie.datasource.write.operation", "upsert") \
    .option("hoodie.upsert.shuffle.parallelism", 200) \
    .mode("append") \
    .save(table_path)

# MOR 表写入(流式场景推荐)
stream_df.writeStream.format("hudi") \
    .option("hoodie.table.name", table_name) \
    .option("hoodie.datasource.write.recordkey.field", "order_id") \
    .option("hoodie.datasource.write.partitionpath.field", "dt") \
    .option("hoodie.datasource.write.precombine.field", "order_time") \
    .option("hoodie.datasource.write.table.type", "MERGE_ON_READ") \
    .option("hoodie.datasource.write.operation", "upsert") \
    .option("checkpointLocation", "s3://checkpoints/hudi-orders") \
    .trigger(processingTime="1 minute") \
    .start(table_path)

Clustering / Compaction 配置:

# 异步 Clustering(写入时自动优化小文件)
hudi_options = {
    "hoodie.table.name": table_name,
    "hoodie.datasource.write.recordkey.field": "order_id",
    "hoodie.datasource.write.partitionpath.field": "dt",
    # Clustering
    "hoodie.clustering.async.enabled": "true",
    "hoodie.clustering.async.max.commits": "4",       # 每 4 次提交触发一次
    "hoodie.clustering.plan.strategy.target.file.max.bytes": "134217728",  # 128MB
    "hoodie.clustering.plan.strategy.small.file.limit": "629145600",       # 600MB
    # MOR Compaction
    "hoodie.compact.inline": "true",
    "hoodie.compact.inline.max.delta.commits": "5",  # 每 5 次 Delta Commit 触发
    "hoodie.compact.inline.trigger.strategy":
        "org.apache.hudi.table.action.compact.strategy.NumCommitsAfterLastCompactionTriggerStrategy",
}

# 离线触发 Clustering
spark.sql(f"""
CALL run_clustering(
    table => '{table_name}',
    path => '{table_path}',
    options => 'hoodie.clustering.plan.strategy.target.file.max.bytes=134217728'
)
""")

# 离线触发 Compaction
spark.sql(f"""
CALL run_compaction(
    table => '{table_name}',
    path => '{table_path}'
)
""")

7.4 Spark + Delta Lake

Delta Lake 核心特性:

特性说明
ACID 事务基于事务日志(_delta_log),并发写入通过乐观并发控制保证一致性
Time Travel通过版本号或时间戳读取历史数据快照
Schema Enforcement写入时严格校验 Schema,类型不匹配直接拒绝
Schema Evolution支持新增列、删除列(Delta 1.2+)
MERGE INTO原生 Upsert 语法,支持复杂条件
OPTIMIZE + Z-ORDER合并小文件并按排序列聚簇,加速点查
Change Data Feed增量读取变更数据(类似 Hudi 增量视图)

Spark 读写 Delta 代码示例:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

spark = SparkSession.builder \
    .appName("DeltaDemo") \
    .config("spark.jars.packages",
            "io.delta:delta-spark_2.12:3.1.0") \
    .config("spark.sql.extensions",
            "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog",
            "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

# 写入 Delta 表
df = spark.range(0, 1000).withColumn("value", col("id") * 10)
df.write.format("delta").mode("overwrite").save("s3://my-bucket/delta/events")

# 读取
spark.read.format("delta").load("s3://my-bucket/delta/events").show()

# SQL 建表
spark.sql("""
CREATE TABLE IF NOT EXISTS delta_events (
    id LONG, value LONG
) USING delta LOCATION 's3://my-bucket/delta/events'
""")

Time Travel:

# 按版本号
df_v2 = spark.read.format("delta") \
    .option("versionAsOf", 2) \
    .load("s3://my-bucket/delta/events")

# 按时间戳
df_ts = spark.read.format("delta") \
    .option("timestampAsOf", "2024-01-15T10:30:00Z") \
    .load("s3://my-bucket/delta/events")

# SQL 语法
spark.sql("SELECT * FROM delta_events VERSION AS OF 2")
spark.sql("SELECT * FROM delta_events TIMESTAMP AS OF '2024-01-15 10:30:00'")

MERGE / OPTIMIZE / ZORDER:

# MERGE INTO(Upsert)
updates = spark.createDataFrame([(1, 999), (1001, 500)], ["id", "value"])
updates.createOrReplaceTempView("updates")

spark.sql("""
MERGE INTO delta_events t
USING updates s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *
""")

# OPTIMIZE:合并小文件
spark.sql("OPTIMIZE delta_events")

# Z-ORDER:按高频过滤列聚簇排序,加速点查
spark.sql("OPTIMIZE delta_events ZORDER BY (id)")

# 清理旧快照(VACUUM,默认保留 7 天)
spark.sql("VACUUM delta_events RETAIN 168 HOURS")

7.5 选型建议

场景推荐格式理由
Databricks 生态 / Spark 重度用户Delta Lake官方原生支持,Z-Order + OPTIMIZE 成熟
多引擎共享(Spark + Flink + Trino)Iceberg引擎中立,社区开放,隐藏分区设计优秀
流式 Upsert / 实时数仓Hudi (MOR)原生流式支持,Compaction/Clustering 完善
Flink CDC 实时入湖 + Spark 分析PaimonFlink 原生流批一体,LSM 高吞吐写入
离线批处理 + 偶尔 UpsertIceberg / DeltaCOW 模式即可满足,运维简单
严格 Schema 管理 + 审计合规Iceberg完整的 Schema/分区演化,快照隔离可审计
已投入 Hadoop 生态 / Hive 兼容Hudi / Iceberg均支持 Hive Catalog,迁移成本低

总结:没有"最好"的数据湖格式,只有"最合适"的选型。建议根据团队的计算引擎、更新频率、查询模式和运维能力综合决策。新项目若以 Spark 为核心且需要多引擎支持,Iceberg 是当前最安全的中立选择。


8. Spark 性能调优

性能调优是 Spark 工程师的核心能力。以下从多个维度系统讲解。

8.1 数据倾斜诊断与解决

数据倾斜是最常见的性能问题,表现为:大部分 Task 很快完成,但少数 Task 执行极慢甚至 OOM。

诊断方法:

  • Spark UI 的 Stages 页面查看 Task 的 Shuffle Read/Write 数据量分布
  • 观察是否有 Task 处理的数据量远超其他 Task
  • 在 Spark SQL 中查看 AQE 倾斜处理日志

解决方案:

方案说明适用场景
加盐(Salting)对热点 Key 加随机前缀,分散到不同 Task聚合类倾斜
广播 Join小表广播,避免 ShuffleJoin 倾斜(小表 < 100MB)
MapJoin / BroadcastHashJoinAQE 自动或手动 hint大小表 Join
拆分热点 Key将热点 Key 单独处理后合并复杂聚合
两阶段聚合先加随机前缀局部聚合,再去前缀全局聚合groupBy 倾斜
AQE Skew JoinSpark 3.0 自动检测并拆分Join 倾斜(推荐)
# 加盐示例:对热点 Key 打散
from pyspark.sql.functions import concat, lit, rand, col

# 第一步:加随机前缀
salted = df.withColumn("salted_key", concat(col("key"), lit("_"), (rand() * 10).cast("int")))
partial = salted.groupBy("salted_key").agg(sum("value").alias("partial_sum"))

# 第二步:去掉前缀再聚合
result = partial.withColumn("key", split(col("salted_key"), "_")[0]) \
                .groupBy("key").agg(sum("partial_sum").alias("total"))

# 广播 Join 提示
from pyspark.sql.functions import broadcast
df_large.join(broadcast(df_small), "key")

# SQL Hint
spark.sql("SELECT /*+ BROADCAST(s) */ * FROM large l JOIN small s ON l.key = s.key")

8.2 Shuffle 调优

# 1. 调整 Shuffle 分区数(默认 200,通常需调整)
spark.conf.set("spark.sql.shuffle.partitions", "200")  # 通用建议:总核数的 2-3 倍

# 2. 启用 AQE 自动合并分区
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")

# 3. Shuffle 压缩
spark.conf.set("spark.shuffle.compress", "true")
spark.conf.set("spark.shuffle.spill.compress", "true")
spark.conf.set("spark.io.compression.codec", "zstd")  # zstd 比 snappy 压缩率更好

# 4. Shuffle 文件缓冲区
spark.conf.set("spark.shuffle.file.buffer", "1MB")        # 默认 32KB
spark.conf.set("spark.reducer.maxSizeInFlight", "96MB")   # 默认 48MB

# 5. RDD API 中使用 coalesce 减少分区
small_df = large_df.coalesce(10)  # 窄依赖,不 Shuffle
# repartition 会 Shuffle
repartitioned = df.repartition(200, col("date"))

8.3 内存调优

# spark-submit 内存参数
--driver-memory 4g \
--executor-memory 8g \
--executor-cores 4 \
--num-executors 20 \
--conf spark.memory.fraction=0.6 \          # Spark 内存占比(默认 0.6)
--conf spark.memory.storageFraction=0.5     # Storage 占 Spark 内存比例(默认 0.5)

Executor 内存规划原则:

  • 每个 Executor 建议 4-8GB,过大易导致 GC 压力
  • 每个 Executor 2-5 个 Core,过多 Core 导致 HDFS I/O 竞争
  • 预留系统内存:spark.memory.fraction 默认为 0.6,即 60% 给 Spark,40% 给用户和其他开销

序列化选择:

序列化方式优点缺点
Java Serialization(默认)兼容好慢、体积大
Kryo Serialization快 10x、体积小需注册类
spark.conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
# Scala: spark.registerKryoClasses(Array(classOf[MyClass]))

缓存策略:

  • 优先使用 MEMORY_AND_DISK_SER 而非 MEMORY_ONLY,避免 OOM
  • DataFrame 建议用 spark.catalog.cacheTable() 或 df.cache(),列式存储更高效
  • 及时 unpersist() 不再使用的缓存

8.4 并行度调优

  • 合理设置分区数:分区太少 → 并行度不足;分区太多 → Task 调度开销大
  • 经验值:分区数 ≈ 总 Executor Core 数的 2-3 倍
  • 读取时控制并行度:
    • HDFS 文件的 InputSplit 数量
    • spark.sql.files.maxPartitionBytes(默认 128MB)
  • AQE 自动调整:Spark 3.0+ 优先使用 AQE

8.5 数据本地性

Spark 倾向于将 Task 调度到数据所在节点,减少网络传输:

级别说明
PROCESS_LOCAL数据在同一 JVM 中(最快)
NODE_LOCAL数据在同一节点
RACK_LOCAL数据在同一机架
ANY数据在任意位置(最慢)

通常不需要手动调整,但如果集群网络延迟高,可适当增大 spark.locality.wait(默认 3s)。

8.6 小文件问题

问题:大量小文件导致 Task 数量过多,调度开销巨大。

# 1. 写入时合并小文件
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128MB")

# 2. 读取时合并小文件
df = spark.read.option("mergeSchema", "true").parquet("hdfs:///data/")

# 3. 使用 Hive 的 CombineHiveInputFormat(对 Hive 表)
spark.conf.set("spark.hadoop.hive.input.format",
    "org.apache.hadoop.hive.ql.io.CombineHiveInputFormat")

# 4. 写入后用 DISTRIBUTE BY 控制输出文件数
df.write.partitionBy("dt").saveAsTable("table")
# 或在 SQL 中
# INSERT OVERWRITE TABLE t PARTITION(dt) SELECT ... DISTRIBUTE BY dt

8.7 列式存储与数据格式

格式类型压缩特点
Text/CSV行式无/弱通用但慢
JSON行式无解析开销大
Parquet列式Snappy/Gzip/ZstdSpark 默认推荐,列裁剪高效
ORC列式Zlib/Snappy/ZstdHive 生态好,ACID 支持
Avro行式SnappySchema 演进好,Kafka 常用
# 写入 Parquet(推荐格式)
df.write.mode("overwrite") \
    .option("compression", "zstd") \
    .parquet("hdfs:///data/output")

# 分区裁剪 + 谓词下推(Parquet/ORC 自动支持)
spark.sql("SELECT name, age FROM users WHERE dt = '2024-01-01' AND age > 25")
# 只读 dt=2024-01-01 分区,且只扫描 name、age 列

8.8 Join 策略选择

Spark 支持多种 Join 实现,理解其原理对性能调优至关重要:

Join 策略触发条件特点
Broadcast Hash Join (BHJ)小表 < spark.sql.autoBroadcastJoinThreshold(默认 10MB)无 Shuffle,最快
Sort-Merge Join (SMJ)大表 Join 大表(默认)需 Shuffle + 排序,稳定但慢
Shuffled Hash Join (SHJ)配置开启且大小表均适合构建哈希表Shuffle 但不排序,比 SMJ 快
Broadcast Nested Loop Join (BNLJ)无等值条件的 Join(如 CROSS JOIN)无 Shuffle,但 O(n×m)
Cartesian Product无条件 Join结果集可能极大,慎用
# 手动指定 Join 策略 Hint
df1.join(broadcast(df2), "key")  # 推荐:广播小表
spark.sql("SELECT /*+ MERGE(l, r) */ * FROM large l JOIN large r ON l.id = r.id")
spark.sql("SELECT /*+ SHUFFLE_HASH(l, r) */ * FROM large l JOIN medium r ON l.id = r.id")

调优建议:

  • 大小表 Join:优先 Broadcast Hash Join,可适当调大 autoBroadcastJoinThreshold(如 100MB)
  • 大表 Join 大表:确保 Join Key 数据类型一致,避免隐式转换导致 Shuffle
  • 多表 Join:注意 Join 顺序,小表尽量先参与
  • AQE 开启后,运行时自动切换为 Broadcast Join

8.9 分区策略

合理的分区设计是性能优化的基础:

# 写入时按字段分区(目录分区)
df.write.partitionBy("dt", "region").parquet("hdfs:///data/sales")
# 生成目录结构:/data/sales/dt=2024-01-01/region=CN/part-xxx.parquet

# 分桶(Bucket):相同 Key 的数据落在同一文件,避免 Join 时 Shuffle
df.write.bucketBy(32, "user_id").sortBy("user_id").saveAsTable("bucketed_users")

# 读取时分区裁剪自动生效
spark.sql("SELECT * FROM sales WHERE dt = '2024-01-01'")  # 只读一个分区目录
策略适用场景注意事项
分区(Partition)高基数的查询过滤字段(如日期)分区数不宜过多(< 10000)
分桶(Bucket)频繁 Join 或聚合的字段桶数需为 2 的幂,与 Join 表一致
聚簇(Clustering)数据湖中自动优化文件布局Iceberg/Hudi/Delta Lake 支持

8.10 Spark UI 实战调试指南

Spark UI(默认端口 4040)是性能调优最重要的工具,能直观展示作业执行细节。以下逐页讲解排查思路。

Jobs 页面:

  • 展示所有 Job 的状态(Succeeded/Failed/Running)、耗时、Stage 数
  • 点击 Job 进入详情,查看 DAG 可视化图:每个方框代表一个 Stage,箭头表示 Shuffle 依赖
  • 任务耗时分布:关注 Duration 列,若某 Job 耗时远超预期,进入对应 Stage 排查
  • 失败任务定位:Failed Jobs 区域直接显示异常堆栈,点击 Failed Stage 查看 Task 失败日志

Stages 页面(核心!):

  • 展示每个 Stage 的 Task 数、输入/输出数据量、Shuffle Read/Write
  • Shuffle Read 列:关注 Min / 25th / Median / 75th / Max,若 Max 远大于 Median(如 Max 是 Median 的 5 倍以上),说明数据倾斜
  • Shuffle Write 列:Map 端写出的数据量,倾斜同样会表现为分布不均
  • Task Time 分布:GC Time 占比过高说明内存压力大;Scheduler Delay 过高说明资源不足
  • 点击 Stage 详情可查看每个 Task 的指标,支持排序

倾斜识别要点:在 Stages 页面查看 Summary Metrics 表,比较 Task Max 与 Median 的 Shuffle Read Size。如果 Max 是 Median 的数倍甚至数十倍,基本可以确认倾斜。

Storage 页面:

  • 展示缓存的 RDD/DataFrame,包括存储级别(Memory/Deserialized 等)、缓存分区数、占用内存/磁盘大小
  • 缓存命中率:Fraction Cached 列显示已缓存分区比例,若偏低说明部分分区未缓存(Executor 丢失或内存不足被淘汰)
  • 分区大小:Size in Memory/ExternalBlockStore 列可判断分区是否均匀
  • 内存 vs 磁盘:若大量数据在 Disk 上,说明内存不足,需调整缓存策略或增大内存

Environment 页面:

  • 展示所有 Spark 配置属性(System Properties、Classpath、Spark Properties)
  • 配置排查:确认提交参数是否生效(如 spark.sql.shuffle.partitions、spark.executor.memory)
  • 可对比运行配置与预期配置,排查"参数没生效"类问题

Executors 页面:

列排查意义
Task Time (Total/GC)GC Time 占 Total 10%+ 说明内存压力大,需调内存或减少缓存
Shuffle Read/Write各 Executor 的 Shuffle 数据量,差异大说明倾斜
Input Size/Records读取数据量分布,不均匀可能是文件大小不均
Storage Memory已用/总缓存内存,满了会溢写磁盘
Failed Tasks非零需查看日志,可能 OOM 或数据异常
Log 链接点击 stdout/stderr 查看 Executor 日志(YARN 模式需通过 RM 代理)

SQL 页面(Spark SQL / DataFrame 作业必看):

  • 展示 SQL 查询的完整执行计划 DAG,节点显示 Scan/Filter/Project/Join/Aggregate/Sort 等算子
  • 点击节点查看详细指标:行数、数据量、耗时、CodeGen 信息
  • Join 识别:节点显示 SortMergeJoin/BroadcastHashJoin/ShuffledHashJoin,确认是否走了预期的 Join 策略
  • Broadcast 识别:若看到 BroadcastExchange 节点,说明小表被广播;若大表被广播可能 OOM
  • AQE Skew Partition:开启 AQE 后,倾斜 Join 会显示 CustomShuffleReader 节点,标注拆分的分区数
  • Scan 节点:查看 number of files read、size of files read、partition filters、data filters,确认分区裁剪和谓词下推是否生效

Structured Streaming 页面:

  • Input Rate:每秒输入数据条数,反映流量波动
  • Processing Rate:每秒处理数据条数,若持续低于 Input Rate 说明处理能力不足
  • Batch Duration:每个微批的处理耗时,若持续接近 Trigger Interval 说明背压严重
  • Watermark 信息:展示当前 Watermark 位置、迟到数据丢弃情况
  • State Operator:状态算子(聚合/Join)的状态行数、内存占用

实战:通过 UI 定位数据倾斜的步骤:

  1. 打开 Stages 页面,找到耗时最长的 Stage,点击 Description 进入详情
  2. 查看 Summary Metrics 表,对比 Shuffle Read 的 Max 和 Median。例如 Median 为 128MB 而 Max 为 8.5GB,确认倾斜
  3. 点击 Tasks 表,按 Shuffle Read Size 降序排列,找到处理最大数据量的 Task
  4. 记录该 Task 的 Locality Level 和 Executor ID,排除节点本地性问题
  5. 返回 SQL 页面,找到对应 Stage 的执行计划,确认是 Join 还是聚合导致倾斜
  6. 若为 Join 倾斜:检查 Join Key 分布,对热点 Key 加盐或开启 AQE Skew Join
  7. 若为聚合倾斜:使用两阶段聚合(加盐局部聚合 + 去盐全局聚合)
  8. 重新提交作业,在 Stages 页面验证 Task 耗时是否均匀

8.11 资源规划与配置参数

Executor 资源规划公式:

推荐配置:
  --num-executors    = ceil(总数据量 / (executor-cores × 单核算力))
  --executor-cores   = 3 ~ 5(建议 4)
  --executor-memory   = 4g ~ 8g(建议 6g)
  --driver-memory     = 2g ~ 4g(复杂作业建议 4g+)

经验法则:
  1. 每个 Executor 分配 4 Core,6GB 内存为黄金配比
  2. YARN 总核数 = num-executors × executor-cores
  3. 预留 10%~20% 集群资源给系统和其他服务
  4. Executor 数量 = 集群总核数 / executor-cores × 0.8

常见踩坑:

问题原因解决方案
executor-memory 过大导致 GC单 Executor 超过 8GB,JVM GC 停顿时间长控制在 4-8GB,增加 Executor 数量而非单 Executor 内存
cores 过多导致 HDFS 并发单 Executor 超过 5 Core,HDFS 客户端并发写导致超时控制在 3-5 Core,HDFS 并发度 = executor-cores × num-executors
Driver OOMcollect() 大量数据或广播大表增大 driver-memory;避免 collect 大数据集
YARN Container 被 Kill物理内存超 Container 限制增大 spark.yarn.executor.memoryOverhead

YARN 模式内存开销:

YARN Container 内存 = spark.executor.memory + spark.executor.memoryOverhead

memoryOverhead 默认值 = max(executor-memory × spark.kubernetes.memoryOverheadFactor, 384MB)
YARN 模式默认 factor = 0.10

示例:
  --executor-memory 6g
  默认 memoryOverhead = max(6g × 0.10, 384MB) ≈ 614MB
  Container 总内存 = 6g + 614MB ≈ 6.7GB

调优建议:
  --conf spark.yarn.executor.memoryOverhead=1g   # 显式指定,避免默认不足

关键配置参数速查表:

参数名默认值说明推荐值
spark.sql.shuffle.partitions200Shuffle 后分区数总核数 × 2-3(开启 AQE 后可适当调大)
spark.executor.memory1g每个 Executor 堆内存4g-8g
spark.executor.cores1每个 Executor CPU 核数4
spark.executor.instances—Executor 数量(YARN/K8s)按集群资源计算
spark.sql.adaptive.enabledfalse(3.x 建议 true)开启 AQE 自适应执行true
spark.sql.adaptive.coalescePartitions.enabledtrueAQE 自动合并小分区true
spark.sql.adaptive.coalescePartitions.minPartitionNum1合并后最小分区数总核数 × 1
spark.sql.adaptive.advisoryPartitionSizeInBytes64MBAQE 目标分区大小128MB
spark.sql.adaptive.skewJoin.enabledtrueAQE 自动处理倾斜 Jointrue
spark.sql.adaptive.skewJoin.skewedPartitionFactor5倾斜分区判定因子(倍数)5(默认即可)
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256MB倾斜分区最小数据量阈值256MB
spark.sql.autoBroadcastJoinThreshold10MB小表自动广播阈值100MB(视内存调整)
spark.serializerorg.apache.spark.serializer.JavaSerializer序列化器org.apache.spark.serializer.KryoSerializer
spark.memory.fraction0.6Spark 内存占 Executor 比例0.6(默认,缓存多可降至 0.5)
spark.memory.storageFraction0.5Storage 占 Spark 内存比例0.5(默认)
spark.speculationfalse推测执行(慢 Task 备份)true(集群空闲时开启)
spark.sql.files.maxPartitionBytes128MB读取文件时单个分区最大字节128MB(大集群可调至 256MB)
spark.sql.files.openCostInBytes4MB打开文件的估算开销(小文件合并用)8MB(小文件多时增大)
spark.sql.broadcastTimeout300s广播等待超时600(大表广播时)
spark.network.timeout120s网络通信超时300s(大 Shuffle 时)
spark.streaming.kafka.consumer.cache.enabledtrue缓存 Kafka 消费者true(默认,避免重复创建)
spark.sql.crossJoin.enabledfalse允许笛卡尔积false(默认,禁止意外笛卡尔积)
spark.shuffle.service.enabledfalseExternal Shuffle Servicetrue(YARN 动态分配必须)
spark.dynamicAllocation.enabledfalse动态资源分配true(配合 ESS)
spark.sql.session.timeZoneJVM 默认会话时区Asia/Shanghai

9. Spark 常见问题与踩坑

9.1 OOM(OutOfMemoryError)

类型原因解决方案
Driver OOMcollect() 拉取大量数据到 Driver用 take()/sample()/write 替代;增大 driver-memory
Executor OOM数据倾斜、缓存过多、大表广播解决倾斜;调整缓存策略;增大 executor-memory
GC OverheadJVM 堆外内存不足增大 spark.executor.memoryOverhead(默认 executor-memory × 0.1)
Python Worker OOMPandas UDF 内存占用大增大 Arrow 批量;分批处理
# 堆外内存配置(重要!)
--conf spark.executor.memoryOverhead=2g \
--conf spark.driver.memoryOverhead=1g \
--conf spark.python.worker.memory=2g

9.2 数据倾斜

见 8.1 节。典型症状:99% 的 Task 很快完成,1% 的 Task 卡住不动。

9.3 Shuffle Fetch Failed

org.apache.spark.shuffle.FetchFailedException:
Failed to connect to /xxx:xxxx

常见原因与解决:

  • Executor 内存不足导致被 YARN Kill:增大 executor-memory / memoryOverhead
  • 网络超时:增大 spark.shuffle.io.maxRetries(默认 3)和 spark.shuffle.io.retryWait(默认 5s)
  • Shuffle 文件过大:增加 Shuffle 分区数
  • 节点故障:开启 spark.shuffle.service.enabled=true(External Shuffle Service)

9.4 序列化错误

org.apache.spark.SparkException: Task not serializable

原因:在 RDD/DataFrame 的闭包中引用了不可序列化的对象(如数据库连接、非序列化的外部类)。

解决:

  • 在函数内部创建不可序列化对象(如数据库连接),不要在 Driver 端创建后传入
  • 使用 foreachPartition / mapPartitions,每个分区创建一次连接
  • 使用 Kryo 序列化并注册类
# ✅ 正确:在分区内创建连接
def process_partition(rows):
    conn = create_db_connection()  # 每个分区创建一次
    for row in rows:
        conn.insert(row)
    conn.close()

df.foreachPartition(process_partition)

9.5 时区问题

# Spark 默认使用 UTC 时区,可能导致时间差 8 小时
spark.conf.set("spark.sql.session.timeZone", "Asia/Shanghai")

# 读取时指定时区
df = spark.read.option("timestampFormat", "yyyy-MM-dd HH:mm:ss") \
    .option("timeZone", "Asia/Shanghai").csv("data.csv")

9.6 UDF 性能问题

  • Python UDF 性能差,优先使用 Spark SQL 内置函数
  • 必须用 UDF 时,使用 Pandas UDF(向量化)
  • 避免在 UDF 中创建重型对象,用 mapPartitions 替代

9.7 其他常见坑

问题说明
count() 触发重算未 cache 的 RDD/DataFrame,每次 Action 都从头计算
collect() 内存溢出确保结果集不超过 Driver 内存
隐式类型转换错误数字与字符串比较时注意类型
Hive 分区不生效执行 MSCK REPAIR TABLE 或添加分区
Spark UI 看不到YARN cluster 模式下通过 ResourceManager 代理访问
文件已存在报错使用 .mode("overwrite") 或先删除

10. 端到端实战项目

10.1 项目背景:电商用户行为日志分析平台

项目目标:构建一个流批一体的电商用户行为分析平台,实时采集用户点击流数据,结合业务库数据,完成实时指标统计和离线报表分析。

数据源:

数据源内容采集方式
Kafka 用户点击流页面浏览、点击、搜索、加购等行为事件Flume/SDK → Kafka
MySQL 订单库订单主表、订单明细Flink CDC / Spark JDBC 全量+增量
MySQL 用户库用户基本信息、地域信息每日全量同步
MySQL 商品库商品分类、价格、品牌每日全量同步

技术栈:

  • 计算引擎:Spark 3.5(Structured Streaming 实时 + Spark SQL 离线)
  • 数据湖:Apache Iceberg(ODS/DWD/DWS/ADS 分层存储)
  • 消息队列:Kafka
  • 缓存:Redis(实时指标查询)
  • 业务库:MySQL(ADS 报表数据)
  • 资源调度:YARN / K8s
  • 监控:Prometheus + Grafana

分层架构:
在这里插入图片描述

10.2 项目架构图

在这里插入图片描述

10.3 核心代码

1. Kafka 消费 + 数据清洗写入 Iceberg ODS:

from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, col, current_timestamp
from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType

spark = SparkSession.builder \
    .appName("UserBehaviorODS") \
    .config("spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.local.type", "hadoop") \
    .config("spark.sql.catalog.local.warehouse", "s3://lakehouse/warehouse") \
    .getOrCreate()

schema = StructType([
    StructField("user_id", LongType()),
    StructField("event_type", StringType()),   # view/click/cart/search/buy
    StructField("product_id", LongType()),
    StructField("category_id", LongType()),
    StructField("event_time", StringType()),
    StructField("device", StringType()),
    StructField("ip", StringType()),
])

kafka_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "kafka1:9092,kafka2:9092") \
    .option("subscribe", "user_behavior") \
    .option("startingOffsets", "latest") \
    .option("maxOffsetsPerTrigger", 100000) \
    .load()

parsed = kafka_df.select(
    from_json(col("value").cast("string"), schema).alias("data")
).select("data.*") \
 .withColumn("ingest_time", current_timestamp()) \
    .withColumn("dt", col("event_time").substr(1, 10))

# 写入 Iceberg ODS(追加模式)
parsed.writeStream \
    .format("iceberg") \
    .option("path", "local.ods.user_behavior") \
    .option("checkpointLocation", "s3://checkpoints/ods_user_behavior") \
    .trigger(processingTime="30 seconds") \
    .outputMode("append") \
    .start()

2. 维度 Join(Broadcast)写入 DWD:

from pyspark.sql.functions import broadcast

# 加载维度表(小表广播)
dim_user = spark.read.format("iceberg").load("local.dim.dim_user")
dim_product = spark.read.format("iceberg").load("local.dim.dim_product")

# 从 ODS 流式读取
ods_stream = spark.readStream \
    .format("iceberg") \
    .option("stream-from-timestamp", str(int((__import__('time').time() - 86400) * 1000))) \
    .load("local.ods.user_behavior")

# 维度关联:广播小表避免 Shuffle
dwd = ods_stream.alias("o") \
    .join(broadcast(dim_user).alias("u"),
          col("o.user_id") == col("u.user_id"), "left") \
    .join(broadcast(dim_product).alias("p"),
          col("o.product_id") == col("p.product_id"), "left") \
    .select(
        col("o.user_id"),
        col("u.age_group"),
        col("u.city"),
        col("o.event_type"),
        col("o.product_id"),
        col("p.category_name"),
        col("p.brand"),
        col("p.price"),
        col("o.event_time"),
        col("o.dt"),
    )

# Upsert 写入 DWD
dwd.writeStream \
    .format("iceberg") \
    .option("path", "local.dwd.dwd_user_behavior") \
    .option("checkpointLocation", "s3://checkpoints/dwd_user_behavior") \
    .trigger(processingTime="1 minute") \
    .start()

3. 窗口聚合写入 DWS:

from pyspark.sql.functions import window, count, sum as _sum, when, col

# 实时窗口聚合:每 1 分钟滚动窗口,统计各分类 PV/UV/加购数
dws_agg = dwd \
    .withWatermark("event_time", "10 minutes") \
    .groupBy(
        window(col("event_time"), "1 minute"),
        col("category_name"),
        col("dt"),
    ).agg(
        count("*").alias("pv"),
        count("user_id").alias("event_count"),
        _sum(when(col("event_type") == "cart", 1).otherwise(0)).alias("cart_count"),
        _sum(when(col("event_type") == "buy", col("price")).otherwise(0)).alias("gmv"),
    ).select(
        col("window.start").alias("window_start"),
        col("window.end").alias("window_end"),
        col("category_name"),
        col("pv"),
        col("event_count"),
        col("cart_count"),
        col("gmv"),
        col("dt"),
    )

dws_agg.writeStream \
    .format("iceberg") \
    .option("path", "local.dws.dws_category_realtime") \
    .option("checkpointLocation", "s3://checkpoints/dws_category_rt") \
    .outputMode("append") \
    .trigger(processingTime="1 minute") \
    .start()

4. ADS 结果输出到 MySQL/Redis:

def write_to_mysql(batch_df, batch_id):
    """将每批次结果写入 MySQL(幂等:REPLACE INTO)"""
    batch_df.write \
        .format("jdbc") \
        .option("url", "jdbc:mysql://mysql:3306/ads") \
        .option("dbtable", "ads_category_realtime") \
        .option("user", "spark") \
        .option("password", "xxx") \
        .option("driver", "com.mysql.cj.jdbc.Driver") \
        .mode("append") \
        .save()

def write_to_redis(batch_df, batch_id):
    """将实时指标写入 Redis,供大屏查询"""
    import redis
    r = redis.Redis(host='redis', port=6379, db=0)
    for row in batch_df.collect():
        key = f"ads:category:{row['category_name']}:{row['window_start']}"
        r.hset(key, mapping={
            "pv": row["pv"],
            "cart_count": row["cart_count"],
            "gmv": float(row["gmv"]),
        })
        r.expire(key, 86400)

# 双流写入
dws_agg.writeStream \
    .foreachBatch(lambda df, id: (write_to_mysql(df, id), write_to_redis(df, id))) \
    .option("checkpointLocation", "s3://checkpoints/ads_sink") \
    .trigger(processingTime="1 minute") \
    .start()

5. 离线补数(批式回刷历史数据):

from datetime import datetime, timedelta

def backfill(start_date, end_date):
    """批量回刷指定日期范围的数据"""
    spark = SparkSession.builder.appName("Backfill").getOrCreate()

    current = datetime.strptime(start_date, "%Y-%m-%d")
    end = datetime.strptime(end_date, "%Y-%m-%d")

    while current <= end:
        dt = current.strftime("%Y-%m-%d")
        print(f"Backfilling {dt} ...")

        # 读取 ODS 指定分区
        ods_df = spark.read.format("iceberg") \
            .load("local.ods.user_behavior") \
            .where(col("dt") == dt)

        # 关联维度
        dwd_df = ods_df.alias("o") \
            .join(broadcast(dim_user), "user_id", "left") \
            .join(broadcast(dim_product), "product_id", "left")

        # 覆盖写入 DWD 指定分区
        dwd_df.writeTo("local.dwd.dwd_user_behavior") \
            .overwritePartitions()

        current += timedelta(days=1)

    spark.stop()

# 执行回刷
backfill("2024-01-01", "2024-01-15")

10.4 项目要点

幂等写入与 Checkpoint 管理:

  • Checkpoint 目录存储 Kafka offset、聚合状态和事务元数据,删除 Checkpoint = 从头消费
  • 生产环境 Checkpoint 应放在可靠存储(HDFS/S3),并设置生命周期管理
  • 幂等写入策略:Iceberg 通过 MERGE INTO 按主键 Upsert;MySQL 使用 REPLACE INTO 或 INSERT ... ON DUPLICATE KEY UPDATE
  • 重新部署时保留 Checkpoint,仅在需要重置时手动删除

延迟数据处理(Watermark + allowedLateness):

# 设置 10 分钟 Watermark:允许事件时间最多比处理时间晚 10 分钟
df.withWatermark("event_time", "10 minutes")

# Iceberg 流式写入可配合 to-snapshot 处理迟到数据
# 超过 Watermark 的数据会被丢弃,但 Iceberg 的 Time Travel 可追溯
# 对于重要的迟到数据,可在 ODS 层保留全量,通过离线补数修复 DWS/DWS
  • Watermark 阈值根据业务延迟特征设置(一般 5-30 分钟)
  • 对于严重延迟的数据(如客户端断网数小时),建议通过离线批处理补数修正
  • 在 ODS 层保留原始数据(Append Only,不删不改),作为"真相源"

监控告警(Streaming Query Listener):

from pyspark.sql.streaming import StreamingQueryListener

class MyListener(StreamingQueryListener):
    def onQueryStarted(self, event):
        print(f"Query started: {event.id}")

    def onQueryProgress(self, event):
        progress = event.progress
        print(f"Batch {progress.batchId}: "
              f"input={progress.numInputRows} rows, "
              f"rate={progress.inputRowsPerSecond:.1f}/s, "
              f"duration={progress.batchDuration}ms")

        # 告警条件:处理速率低于输入速率(消费滞后)
        if progress.inputRowsPerSecond > progress.processedRowsPerSecond * 1.5:
            send_alert(f"Consumer lag! input={progress.inputRowsPerSecond}, "
                       f"processed={progress.processedRowsPerSecond}")

        # 告警:批次耗时异常
        if progress.batchDuration > 120000:  # 超过 2 分钟
            send_alert(f"Batch duration too long: {progress.batchDuration}ms")

    def onQueryTerminated(self, event):
        print(f"Query terminated: {event.id}, exception={event.exception}")
        if event.exception:
            send_alert(f"Streaming query failed: {event.exception}")

spark.streams.addListener(MyListener())

def send_alert(msg):
    """对接 Prometheus AlertManager / 钉钉 / 企业微信"""
    import requests
    requests.post("https://alert-webhook.example.com/", json={"text": msg})

常见问题:

问题原因解决方案
小文件过多每批次产生大量小 Parquet 文件开启 Iceberg write.distribution-mode=hash;定期执行 rewrite_data_files 合并
Kafka offset 丢失Checkpoint 损坏或被误删备份 Checkpoint;设置 kafka.group.id 配合外部 Offset 管理
Schema 变更上游新增字段导致流解析失败使用 from_json 时开启 mergeSchema=true;Iceberg 支持 Schema 演化新增列
状态过大无限聚合导致 State 膨胀设置 Watermark 清理过期状态;使用 RocksDB State Backend;定期 Rebalance
Executor OOM大窗口聚合或数据倾斜增大内存;加盐打散热点 Key;开启 AQE Skew Join
写入冲突并发写同一 Iceberg 表Iceberg 支持乐观并发,重试事务;避免流和批同时写同一分区
Kafka 消费积压处理能力不足增加 Executor/分区数;优化处理逻辑;使用 maxOffsetsPerTrigger 限流

11. Spark 3.x 新特性

11.1 AQE(自适应查询执行)

Spark 3.0 最重要的特性,详见 5.8 节。Spark 3.x 持续增强:

  • Spark 3.2:AQE 支持 BroadcastNestedLoopJoin 优化
  • Spark 3.5:AQE 对窗口函数的优化增强

11.2 Dynamic Partition Pruning(动态分区裁剪)

Spark 3.0 引入,在星型模型的 Join 中,自动根据维度表过滤条件裁剪事实表的分区,避免扫描无关数据。

-- 维度表 products 只查询了 category='Electronics'
-- DPP 会自动将事实表 sales 的分区裁剪到相关产品
SELECT * FROM sales s
JOIN products p ON s.product_id = p.id
WHERE p.category = 'Electronics';

11.3 Pandas API on Spark

Spark 3.2 引入(原 Koalas 项目合并),让 Pandas 用户可以几乎无学习成本地使用 Spark:

import pyspark.pandas as ps

# Pandas 风格的 API,自动分布式执行
pdf = ps.read_csv("hdfs:///data/huge_file.csv")
result = pdf.groupby("category")["amount"].sum().reset_index()
result.to_parquet("hdfs:///output/summary.parquet")

# Pandas API on Spark 与 PySpark DataFrame 互转
sdf = result.to_spark()
pdf = sdf.to_pandas_on_spark()

11.4 Spark Connect

Spark 3.4 正式发布,将客户端与 Spark 集群解耦。客户端(瘦客户端)通过 gRPC 与 Spark Driver 通信,支持:

  • 任意 IDE / Notebook 作为客户端,无需在本地安装 Spark
  • 多语言、多版本客户端连接同一 Spark 集群
  • 微服务架构中嵌入 Spark 能力
# 使用 Spark Connect 连接远程集群
spark = SparkSession.builder \
    .remote("sc://spark-driver:15002") \
    .getOrCreate()

11.5 ANSI SQL 模式

Spark 3.0 开始引入 ANSI SQL 兼容模式,Spark 3.4/3.5 大幅完善:

spark.conf.set("spark.sql.ansi.enabled", "true")

启用后:

  • 除零报错(而非返回 NULL)
  • 整数溢出报错
  • 字符串转数字失败报错
  • 类型转换更严格
  • 保留关键字与标准 SQL 一致

11.6 Python 性能提升

  • PySpark 优化:Spark 3.x 大幅改善了 Python API 的性能
  • Arrow 优化:Pandas UDF / Pandas API on Spark 使用 Apache Arrow 实现高效数据交换
  • Python 进程复用:Python Worker 复用,减少启动开销
  • Spark 3.5:Python UDF 支持 profiling,便于性能分析
  • Spark 4.0 预览:进一步优化 Python 路径,引入更高效的 DataFrame 内部表示

11.7 其他重要特性

版本特性
3.0AQE、DPP、加速器感知调度(GPU/YARN)、Java 11+
3.1Predicate Hint、Session Window、ANSI 间隔类型
3.2Pandas API on Spark、Session Window、ANSI 模式增强
3.3Row-level runtime filtering、TIMESTAMP_NTZ 类型
3.4Spark Connect 稳定、Python Profiler、默认 UTC 时区
3.5Python UDT、Named Variable、Arrow 优化、Lateral Column Alias
4.0(预览)更强 ANSI SQL、列式处理优化、更好的 Python 体验、Variant 类型

12. 学习路线与实战建议

12.1 学习阶段划分

在这里插入图片描述

12.2 推荐资源

书籍:

  • 《Spark 权威指南》(Bill Chambers 等)—— 最适合入门
  • 《Spark 快速大数据分析》(Holden Karau 等)—— RDD 经典
  • 《Spark SQL 内核剖析》—— 深入 Catalyst 源码
  • 《高性能 Spark》—— 性能调优进阶

官方文档:

在线课程:

  • Databricks Academy(官方培训)
  • edX / Coursera 上的 Berkeley 大数据课程

源码与社区:

  • GitHub: https://github.com/apache/spark
  • Stack Overflow 的 apache-spark 标签
  • Databricks Blog(最佳实践和技术深度文章)

12.3 实战项目建议

项目技术点难度
网站日志分析RDD/DataFrame、ETL、聚合⭐
电商用户行为分析Spark SQL、窗口函数、漏斗分析⭐⭐
实时风控系统Structured Streaming、Kafka、CEP⭐⭐⭐
推荐系统特征工程MLlib、特征拼接、模型训练⭐⭐⭐
数据湖 ETL 管道Iceberg/Hudi、分区演进、流批一体⭐⭐⭐⭐
Spark on K8s 部署实践Docker 镜像、K8s 配置、弹性调度⭐⭐⭐⭐
电商实时数仓Kafka + SS + Iceberg + Redis、Watermark⭐⭐⭐⭐⭐
大规模数据倾斜治理AQE、加盐、Broadcast Join、Spark UI 排查⭐⭐⭐⭐

12.4 面试高频考点

  1. Spark 架构:Driver/Executor 工作流程,YARN cluster vs client 区别
  2. RDD 原理:五大特性、宽窄依赖、血缘容错
  3. Shuffle:Shuffle 过程、SortShuffle 与 HashShuffle 区别、Shuffle 调优
  4. 数据倾斜:成因、诊断、至少三种解决方案
  5. Spark SQL:Catalyst 优化流程、AQE 原理、Tungsten
  6. 内存管理:统一内存模型、各区域比例、OOM 排查
  7. 持久化:cache/persist/checkpoint 区别
  8. Join 策略:BroadcastHashJoin、SortMergeJoin、ShuffledHashJoin
  9. Exactly-Once:Structured Streaming 如何保证
  10. Spark 3.x 新特性:AQE、DPP、Spark Connect
  11. 数据湖选型:Delta Lake / Iceberg / Hudi 对比,各自适用场景
  12. Spark UI 倾斜排查:如何通过 Stages 页面的 Task Max/Median 定位数据倾斜
  13. Iceberg Time Travel 原理:Snapshot/Manifest 机制,如何实现历史版本读取
  14. Hudi COW vs MOR:写入/读取差异、Compaction 时机、选型建议
  15. Spark on K8s vs YARN:架构差异、弹性/隔离/生态对比、适用场景
  16. Iceberg 分区演化:与 Hive 分区的区别,如何实现隐式分区
  17. Hudi Compaction/Clustering:作用、触发策略、对读写性能的影响

总结

Spark 作为大数据领域最主流的计算引擎,其生态完整、社区活跃、就业需求旺盛。学习 Spark 的关键路径是:

  1. 先建立全局认知:理解架构和运行原理,不要只停留在 API 调用
  2. 动手实践:Local 模式搭起来,写代码跑通比看十篇文章都有用
  3. 读懂 Spark UI:它是诊断性能问题的最佳工具
  4. 深入 SQL:Spark SQL 是当前和未来的核心,DataFrame API 和 SQL 都要熟练
  5. 关注数据倾斜:这是生产环境最常遇到的问题
  6. 与时俱进:关注 AQE、Spark Connect、Pandas API on Spark 等新特性,以及 Spark 4.0 的方向

大数据技术栈虽广,但 Spark 是那个"牵一发而动全身"的核心引擎。掌握了 Spark,你就拿到了数据工程领域的一把关键钥匙。

希望这份指南能帮你在 Spark 学习路上少走弯路。祝你学习愉快,早日成为 Spark 高手!🚀

更多推荐