Spark 从零到进阶:数据开发工程师的系统学习指南
写在前面:如果你是一名刚接触 Spark 的数据开发工程师,面对网上零散的教程和概念感到无从下手,那么这篇文章就是为你准备的。我们将从"Spark 是什么"出发,一路走到性能调优与生产踩坑,配合大量代码示例和架构图,帮你建立完整的知识体系。本文基于 Spark 3.5.x 版本编写,并会提及 Spark 4.0 的预览方向。
目录
1. Spark 概述
1.1 什么是 Spark
Apache Spark 是一个开源的统一分析引擎,专为大规模数据处理而设计。它最初于 2009 年在加州大学伯克利分校 AMPLab 诞生,2010 年开源,2014 年成为 Apache 顶级项目。Spark 提供了 SQL、流计算、机器学习和图计算等一整套大数据处理能力,可以运行在 Hadoop YARN、Kubernetes、Standalone 等多种集群管理器上。
1.2 发展历史
| 2009 | Spark 在 UC Berkeley AMPLab 诞生 |
| 2010 | BSD 许可开源 |
| 2013 | 捐赠给 Apache 软件基金会 |
| 2014 | Apache 顶级项目;Spark 1.0 发布 |
| 2016 | Spark 2.0:DataFrame/Dataset 成为主流 API,Structured Streaming |
| 2020 | Spark 3.0:AQE、Dynamic Partition Pruning |
| 2023 | Spark 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 依然是重要的存储和资源管理组件。
| 中间结果 | 落盘(HDFS) | 优先内存,可溢写磁盘 |
| 编程模型 | Map + Reduce 两阶段 | DAG(有向无环图),多阶段 |
| 延迟 | 高(批处理) | 低(批流统一) |
| 迭代计算 | 每次迭代读写 HDFS | 内存缓存,迭代效率高 10-100x |
| 生态 | 仅 MapReduce | SQL/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 | 单机运行,多线程模拟分布式 | 开发调试、学习 |
| Standalone | Spark 自带集群管理器 | 小规模独立集群 |
| 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 多个分区使用,需要 Shuffle | reduceByKey、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-shell | Scala 交互式 REPL,适合探索和调试 |
| pyspark | Python 交互式 REPL |
| spark-sql | SQL 交互式命令行 |
| spark-submit | 提交应用到集群(生产环境) |
spark-submit 常用参数:
| –master | 集群管理器 URL | yarn / spark://host:7077 / local[*] |
| –deploy-mode | Driver 运行位置 | cluster / client |
| –name | 应用名称 | MySparkApp |
| –class | Java/Scala 主类 | com.example.MyApp |
| –jars | 额外依赖 JAR | –jars lib/mysql-connector.jar |
| –packages | Maven 依赖自动下载 | –packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1 |
| –files | 分发文件到工作目录 | –files config.properties |
| –conf | Spark 配置项 | –conf spark.sql.shuffle.partitions=200 |
| –driver-memory | Driver 内存 | 4g |
| –executor-memory | 每个 Executor 内存 | 8g |
| –executor-cores | 每个 Executor CPU 核数 | 4 |
| –num-executors | Executor 数量(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 模式对比:
| 部署方式 | 依赖 Hadoop YARN 集群 | 容器化部署,镜像打包所有依赖 |
| 弹性扩缩容 | 依赖 YARN 队列配置,扩容较慢 | 秒级 Pod 创建,弹性能力强 |
| 资源隔离 | 队列级隔离,Container 共享 OS | Pod 级隔离,容器边界清晰 |
| 依赖管理 | 通过 –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.image | Docker 镜像地址 | registry/spark:3.5.1 |
| spark.kubernetes.namespace | K8s 命名空间 | spark-jobs |
| spark.kubernetes.authenticate.driver.serviceAccountName | Driver 使用的 ServiceAccount | spark |
| spark.kubernetes.driver.pod.name | Driver Pod 名称 | spark-app-driver |
| spark.kubernetes.executor.podNamePrefix | Executor Pod 名称前缀 | spark-app-exec |
| spark.kubernetes.driver.limit.cores | Driver CPU 限制 | 2 |
| spark.kubernetes.executor.limit.cores | Executor CPU 限制 | 4 |
| spark.kubernetes.driver.request.cores | Driver CPU 请求 | 1 |
| spark.kubernetes.executor.request.cores | Executor CPU 请求 | 2 |
| spark.kubernetes.memoryOverheadFactor | 堆外内存比例 | 0.2 |
| spark.kubernetes.allocation.batch.size | 每批申请 Pod 数 | 10 |
| spark.kubernetes.executor.deleteOnTermination | 完成后删除 Executor Pod | true |
| spark.kubernetes.driver.podTemplateFile | Driver Pod 模板文件 | driver-pod.yaml |
| spark.kubernetes.executor.podTemplateFile | Executor Pod 模板文件 | executor-pod.yaml |
| spark.kubernetes.file.upload.path | 应用依赖上传路径 | s3://spark-uploads/ |
| spark.kubernetes.authenticate.caCertFile | CA 证书路径 | /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 五大特性
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 触发真正的计算。
| Transformation | map、filter、flatMap、mapPartitions、sample、union、intersection、distinct、groupByKey、reduceByKey、aggregateByKey、sortByKey、join、cogroup、cartesian、coalesce、repartition |
| Action | collect、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
这是面试高频考点:
| 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 | 内存/磁盘可选 | 保留 | 根据数据量选择 |
| checkpoint | HDFS | 切断 | 长血缘链、迭代计算 |
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
| 数据模型 | 无结构 | 带 Schema 的行 | 带 Schema 的强类型对象 |
| 类型安全 | 是(编译时) | 否(运行时) | 是(编译时) |
| 优化 | 无 | Catalyst | Catalyst |
| 语言 | 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 效率和内存效率:
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
| Kafka | 最常用,支持从 Kafka 读取消息 |
| File | 监听目录中新文件 |
| Socket | 测试用,从 TCP Socket 读取 |
| Rate | 测试用,每秒生成指定行数 |
| Kafka | 写入 Kafka Topic |
| File | 写入文件(Parquet/JSON/CSV) |
| Console | 控制台(调试用) |
| Foreach/ForeachBatch | 自定义写入逻辑 |
| Memory | 存储为内存表(调试用) |
6.5 Exactly-Once 语义
Structured Streaming 通过以下机制保证精确一次(Exactly-Once):
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 语义,兼具数据湖的灵活性和数据仓库的管理能力。
| 定位 | 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/Hive | Spark/Flink/Trino/Presto |
| 社区活跃度 | Databricks 商业驱动 | Apache 顶级项目,中立开放 | Apache 顶级项目,Uber 起源 |
| 典型用户 | Databricks 客户 | Apple/Netflix/LinkedIn | Uber/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 表类型对比:
| 写入方式 | 每次写入更新重写整个 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 文件 |
| Compaction | MOR 表将 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 分析 | Paimon | Flink 原生流批一体,LSM 高吞吐写入 |
| 离线批处理 + 偶尔 Upsert | Iceberg / Delta | COW 模式即可满足,运维简单 |
| 严格 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 | 小表广播,避免 Shuffle | Join 倾斜(小表 < 100MB) |
| MapJoin / BroadcastHashJoin | AQE 自动或手动 hint | 大小表 Join |
| 拆分热点 Key | 将热点 Key 单独处理后合并 | 复杂聚合 |
| 两阶段聚合 | 先加随机前缀局部聚合,再去前缀全局聚合 | groupBy 倾斜 |
| AQE Skew Join | Spark 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/Zstd | Spark 默认推荐,列裁剪高效 |
| ORC | 列式 | Zlib/Snappy/Zstd | Hive 生态好,ACID 支持 |
| Avro | 行式 | Snappy | Schema 演进好,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 实现,理解其原理对性能调优至关重要:
| 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 定位数据倾斜的步骤:
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 OOM | collect() 大量数据或广播大表 | 增大 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.partitions | 200 | Shuffle 后分区数 | 总核数 × 2-3(开启 AQE 后可适当调大) |
| spark.executor.memory | 1g | 每个 Executor 堆内存 | 4g-8g |
| spark.executor.cores | 1 | 每个 Executor CPU 核数 | 4 |
| spark.executor.instances | — | Executor 数量(YARN/K8s) | 按集群资源计算 |
| spark.sql.adaptive.enabled | false(3.x 建议 true) | 开启 AQE 自适应执行 | true |
| spark.sql.adaptive.coalescePartitions.enabled | true | AQE 自动合并小分区 | true |
| spark.sql.adaptive.coalescePartitions.minPartitionNum | 1 | 合并后最小分区数 | 总核数 × 1 |
| spark.sql.adaptive.advisoryPartitionSizeInBytes | 64MB | AQE 目标分区大小 | 128MB |
| spark.sql.adaptive.skewJoin.enabled | true | AQE 自动处理倾斜 Join | true |
| spark.sql.adaptive.skewJoin.skewedPartitionFactor | 5 | 倾斜分区判定因子(倍数) | 5(默认即可) |
| spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes | 256MB | 倾斜分区最小数据量阈值 | 256MB |
| spark.sql.autoBroadcastJoinThreshold | 10MB | 小表自动广播阈值 | 100MB(视内存调整) |
| spark.serializer | org.apache.spark.serializer.JavaSerializer | 序列化器 | org.apache.spark.serializer.KryoSerializer |
| spark.memory.fraction | 0.6 | Spark 内存占 Executor 比例 | 0.6(默认,缓存多可降至 0.5) |
| spark.memory.storageFraction | 0.5 | Storage 占 Spark 内存比例 | 0.5(默认) |
| spark.speculation | false | 推测执行(慢 Task 备份) | true(集群空闲时开启) |
| spark.sql.files.maxPartitionBytes | 128MB | 读取文件时单个分区最大字节 | 128MB(大集群可调至 256MB) |
| spark.sql.files.openCostInBytes | 4MB | 打开文件的估算开销(小文件合并用) | 8MB(小文件多时增大) |
| spark.sql.broadcastTimeout | 300s | 广播等待超时 | 600(大表广播时) |
| spark.network.timeout | 120s | 网络通信超时 | 300s(大 Shuffle 时) |
| spark.streaming.kafka.consumer.cache.enabled | true | 缓存 Kafka 消费者 | true(默认,避免重复创建) |
| spark.sql.crossJoin.enabled | false | 允许笛卡尔积 | false(默认,禁止意外笛卡尔积) |
| spark.shuffle.service.enabled | false | External Shuffle Service | true(YARN 动态分配必须) |
| spark.dynamicAllocation.enabled | false | 动态资源分配 | true(配合 ESS) |
| spark.sql.session.timeZone | JVM 默认 | 会话时区 | Asia/Shanghai |
9. Spark 常见问题与踩坑
9.1 OOM(OutOfMemoryError)
| Driver OOM | collect() 拉取大量数据到 Driver | 用 take()/sample()/write 替代;增大 driver-memory |
| Executor OOM | 数据倾斜、缓存过多、大表广播 | 解决倾斜;调整缓存策略;增大 executor-memory |
| GC Overhead | JVM 堆外内存不足 | 增大 spark.executor.memoryOverhead(默认 executor-memory × 0.1) |
| Python Worker OOM | Pandas 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_da
未经允许不得转载:171主机测评 » Spark 从筑基到化神
相关推荐评论 抢沙发 |