跳至正文
数据工程 — Spark

Spark

术语修订(2026-09-07,Agent:Codex):通过 Context7 和 Spark 官方资料统一 Driver、Executor、Transformation、Action 与 Dataset 等术语,并核对相邻执行语义。此次运行记录:模型 gpt-6-astra,reasoning effort ultra,执行入口 Codex Desktop,提供方 openai,CLI 版本 0.153.4(不代表桌面 App 版本)。

Spark 是内存密集型的分布式大数据计算引擎。Apache Spark™ is a multi-language engine for executing data engineering, data science, and machine learning on single-node machines or clusters1.

Spark 生态圈以 Spark Core 为核心,支持从 HDFS、Amazon S3、HBase、ElasticSearch、MongoDB、MySQL、Kafka 等多种数据源读取数据。同时,当前 Spark 的 Cluster Manager 类型包括 Standalone、Hadoop YARN 与 Kubernetes,负责为应用分配资源,从而完成 Spark 应用程序的计算任务。

搭建环境

首先,去官网下载 Spark,下载下来的 Spark 文件中,常见的目录有 bin、 sbin、kubernetes、data、examples,其中 bin 里面有 pysparkspark-submitspark-sql 等常用的命令。以 Python 为例,启动 ./bin/pyspark 后,可以在 Python Shell 中执行一些 Spark 操作, 通过访问 http://localhost:4040 可以浏览 Spark UI 界面。

相关概念

使用之前还需要了解一些基本概念。

Driver 与 Executor

The driver coordinates a Spark application; its executors run tasks and retain application data on worker nodes. Cluster Mode Overview

SparkSession

可以存在多个 SparkSession 用于连接多个不同的数据库。

执行计划

应用

作业

RDD

DataFrame

Job

一次 Action 可触发 Job,Job 再划分为多个 Stage,每个 Stage 包含一组 Task。Jobs 与 Stages 标签页对应这些执行层次,可继续查看单个 Task 的详情。我们可以看到各粒度下的完成状态、I/O 相关指标、内存消耗、执行持续时间等信息。

Stages

通过 Stages 标签页可以概览各 Job 中 Stage 的当前状态。

Executors

Executors 标签页提供应用各 Executor 的运行信息。

Transformations 与 Actions

Dataset operations fall into two categories: transformations build new Datasets, while actions trigger computation and return or write results. Dataset API

Transformations 采用惰性求值。select()filter() 返回新的 DataFrame,记录需要执行的转换,不会修改原 DataFrame,也不会立即计算全部结果。对 RDD,这种依赖记录通常称为 lineage;对 Dataset / DataFrame,应从 logical plan 理解后续优化与执行。

调用 Action 时,Spark 才执行获得结果所需的计算;读取缓存等情况可能复用已有结果。

Narrow transformations 与 Wide transformations

可以按分区依赖理解 narrow transformations 与 wide transformations,关键是是否需要跨分区交换数据。

Narrow transformations 可在输入分区内完成计算。普通列选择 select() 和行过滤 filter() 不要求跨分区交换数据。

Wide transformations 涉及跨分区依赖,通常需要 shuffle,例如按键聚合或全局排序。是否出现实际的数据交换,应查看 physical plan,不能只凭方法名判断。

DAG

DAG(Directed Acyclic Graph,有向无环图)表达计算依赖。Spark 根据依赖组织 Stage,再将各 Stage 中的 Task 交给 Executor;Executor 可以使用内存或磁盘存储应用数据。Cluster Mode Overview

RDD

RDD(Resilient Distributed Dataset,弹性分布式数据集) 是 Spark 1.0 时期的底层数据对象。

API

Spark SQL

pyspark.sql.dataframe.DataFrame

Pandas API on Spark

pyspark.pandas.frame.DataFrame

Dataset

Dataset 也是 Spark 1.6 加入的 feature。

A DataFrame is a Dataset organized into named columns,  DataFrame is represented by a Dataset of Rows. The Dataset API is available in Scala and Java. 并没有针对 Python 提供 Dataset 接口。

Spark 2.0 用相似的接口统一了 DataFrame API 和 Dataset API。

DataFrame 没法做到强类型检查,而编译型语言 Scala 和 Java 可以,所以 Dataset 相对于 DataFrame 提供了强类型支持。

SQL、DataFrame 和 Dataset 区别

DataFrame 和 Dataset 都是基于 RDD 构建的

https://cdn.yindongliang.com/2023/biJN7A.png

Spark 3.0 中有两种 DataFrame。PySpark 中默认的 DataFrame 其实是 Spark SQL 的 pyspark.sql.dataframe.DataFrame,这个 DataFrame 可以通过接口 pyspark.sql.DataFrame.pandas_api 转换为 pyspark.pandas.frame.DataFrame

底层引擎

上层的 DataFrame API 和 Dataset API 都是由底层的 Spark SQL 引擎支撑的,Spark SQL 引擎的核心是 Catalyst 优化器和 Tungsten 项目。两者共同支撑着高层的 DataFrame API 和 Dataset API,以及 SQL 查询。

优化过程都是一样的:先是构建逻辑计划,接着是生成物理计划,最后是生成紧凑的二进制代码。

执行计划

Spark SQL

在这里讨论的 Spark SQL 文件表场景中,需区分 managed table 与 external table:前者由 Spark 管理表数据的位置和生命周期,后者使用显式指定的外部数据路径。删除使用自定义 path 的表时,Spark 保留该路径中的数据;未指定路径、使用默认 warehouse 位置的表,删除时会一并移除默认数据目录。Saving to Persistent Tables

python
us_flights_df = spark.sql("SELECT * FROM us_delay_flights_tbl")
us_flights_df2 = spark.table("us_delay_flights_tbl")

repartition

部署模式

Spark Standalone

Apache Mesos:历史部署选项,不在当前官方支持的 Cluster Manager 列表中。

YARN

Kubernetes

命令行

spark-shell

spark-submit

其他支持

Structured Streaming 流数据处理

MLlib 机器学习库

参考文档

Spark 2.2.x 中文官方参考文档
https://spark-reference-doc-cn.readthedocs.io/zh_CN/latest/index.html

SparkBy{Examples}
https://sparkbyexamples.com/pyspark-tutorial/

Spark 编程指南
https://doc.yonyoucloud.com/doc/spark-programming-guide-zh-cn/index.html

Spark-Programming-In-Python

Spark 快速大数据分析(第 2 版)

图解Spark 大数据快速分析实战

一些代码示例:

Footnotes

  1. Spark 官网 https://spark.apache.org

本文共 1355 字,创建于 Oct 16, 2023

相关标签:Python

评论

博客助手

正在打开博客助手…