欢迎光临
我们一直在努力

Spark 从筑基到化神

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

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


目录

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

  • 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 依然是重要的存储和资源管理组件。

    对比维度MapReduceSpark
    中间结果 落盘(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 模式对比:

    对比维度Spark on YARNSpark on K8s
    部署方式 依赖 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 五大特性

  • 分区列表(Partitions):数据被切分为多个分区,每个分区在一个节点上计算
  • 计算函数(Compute):每个分区都有一个计算函数来生成数据
  • 依赖关系(Dependencies):RDD 之间有血缘关系,用于故障恢复
  • 分区器(Partitioner):KV 类型 RDD 可选(Hash/Range),决定数据分布
  • 优先位置(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 触发真正的计算。

    类型常见算子
    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

    这是面试高频考点:

    对比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 内存/磁盘可选 保留 根据数据量选择
    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

    特性RDDDataFrameDataset
    数据模型 无结构 带 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 效率和内存效率:

  • 内存管理:使用堆外内存(off-heap),避免 JVM GC 开销
  • 二进制处理:数据以二进制格式存储,避免 Java 对象的序列化/反序列化
  • Whole-Stage CodeGen:将整个 Stage 的多个算子融合为一个 Java 函数,消除虚函数调用
  • 向量化读取: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):

  • 可重放的 Source:如 Kafka,记录 offset 可重新读取
  • 幂等的 Sink:或使用事务写入(如 Kafka 事务、文件原子写入)
  • 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/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 表类型对比:

    特性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 文件
    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 实现,理解其原理对性能调优至关重要:

    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 定位数据倾斜的步骤:

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

    赞(0)
    未经允许不得转载:171主机测评 » Spark 从筑基到化神
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址

    © 2010-2026   171主机测评   网站地图

    请求次数:57 次,加载用时:3.715 秒,内存占用:38.27 MB