什么是 Spark DataFrame API?

Apache Spark 是一个分布式处理引擎,旨在大规模处理结构化、半结构化和非结构化数据。为了与引擎交互并定义分布式数据操作,开发者会使用 Spark API。


现代 Spark 开发的核心是 Spark DataFrame API。DataFrame 以具名列的形式整理和处理分布式数据,结构类似于关系型数据库中的表或电子表格。它提供了一个结构化的声明式接口,用于分析海量数据集,同时在后台自动优化物理执行。

Spark API 架构和执行

如需了解 Spark API 如何在分布式网络中执行代码,必须先了解其运行时架构的基本组件:


  • 驱动程序:Spark 应用的中央协调器。它会读取您的代码,将声明式操作转换为逻辑执行计划,并在工作器节点之间调度任务。
  • 集群管理器:资源分配器(如 Kubernetes 或独立资源管理器),负责在整个集群中预配 CPU、内存和网络资源。
  • 执行器:在集群节点上运行的工作器实例。它们从驱动程序接收任务指令,在内存中本地执行数据处理操作,并返回结果或将其写入存储空间。
  • DAG、阶段和任务:Spark 不会逐行执行数据转换。而是由驱动程序将 DataFrame 代码编译为物理操作符的有向无环图 (DAG)。驱动程序会将这些操作符分组到称为“阶段”的广泛执行阶段(由数据重排边界划分),并将这些阶段分解为分配给执行器并行执行的各个任务。

Spark RDD 与 DataFrame 有何区别?

在设计分布式数据流水线时,开发者需要在两种主要编程接口之间做出选择:


  • Spark RDD(弹性分布式数据集):原始的低级层 Spark API。它表示分布在集群中的不可变、可容错的 JVM 对象集合。由于 RDD 包含任意 Java 对象,因此执行引擎无法检查其内部结构。这迫使开发者手动编写和调优低级层转换逻辑,往往导致大量的垃圾回收开销和缓慢的序列化。
  • Spark DataFrame:结构化数据处理的现代标准。DataFrame 可将数据整理成由严格架构定义的行和列。由于执行引擎了解数据集的数据类型和结构,因此可以在执行代码之前自动优化代码。

为什么选择 DataFrame 而不是 RDD?

对于几乎所有数据工程和数据科学应用场景,DataFrame API 都是首选。RDD 仅适用于必须使用自定义 Java 对象处理原始非结构化数据(如二进制文件或媒体流)的场景。DataFrame 需要的代码更少,自动执行速度更快,并利用堆外二进制内存管理来绕过 JVM 垃圾回收瓶颈。

数据团队如何使用 Spark API

Spark DataFrame API 简化了处理、清理和准备大量数据的计算密集型任务。


  • 数据工程师:工程师使用 DataFrame API 构建弹性、可容错的数据流水线。他们提取原始数据、应用架构、运行结构化转换,并将符合规范的表写入数据湖或云数据仓库。声明式语法使他们能够每天处理数 TB 的数据,同时最大限度地降低代码复杂性和维护开销。
  • 数据科学家:数据科学家使用 DataFrame API 探索和准备超出单台机器内存限制的海量数据集。通过将数据分区分布在工作器节点上,他们可以大规模运行探索性数据分析、清理 null 值并设计特征,从而加快获取分析洞见的速度。

基于 DataFrame API 构建的核心库

Spark DataFrame API 是 Spark 高级处理库的编程基础:


  • Spark SQL:此模块允许您直接对 Spark 数据集执行 ANSI SQL 查询。您可以将 DataFrame 作为临时视图进行查询,将声明式 Python 或 Scala 代码与标准 SQL 查询无缝融合。
  • Structured Streaming:此引擎用于处理连续的实时数据流。它使用与静态批处理相同的 DataFrame API 命令,自动处理微批处理、容错和状态管理。
  • MLlib:Spark 的分布式机器学习库,使用 DataFrame 管理和准备训练数据集,提供内置的可伸缩算法,用于分类、回归、聚类和协同过滤。

Spark DataFrame API 的优势

声明式优化

DataFrame 使用内置的查询优化器,称为 Catalyst 优化器。当您编写 DataFrame 代码时,优化器会自动重写物理执行计划,应用过滤条件下推、投影剪枝和联接重新排序,从而在无需手动调优的情况下尽可能快地运行。

语言灵活性

无论您的团队使用 Python (PySpark)、Scala、Java 还是 R 编写代码,DataFrame API 都能确保相同的性能。底层执行计划在同一个独立于 JVM 的执行层中进行编译和优化。

高效的内存管理

DataFrame 利用 Tungsten 执行引擎,以高度压缩的堆外二进制格式存储数据。这消除了 JVM 对象创建开销,并防止垃圾回收暂停冻结执行器线程。

如何在 Google Cloud 上使用 Spark DataFrame API

传统上,运行开源 Spark 需要手动配置集群、管理软件版本以及进行复杂的网络对等互连。Google Cloud 通过提供 Managed Service for Apache Spark 简化了这一过程,将分布式执行转换为全托管式、企业级数据平台。


数据团队使用以下功能在 Google Cloud 上执行 Spark DataFrame 工作负载:


  • 无服务器和托管式集群部署:Google Cloud 可让您选择适合您运营需求的执行模型。您可以在无服务器模式下运行 DataFrame 代码来直接提交批量作业,只需为资源自动扩缩时的确切运行时间(秒)付费;或者部署高度可自定义的持久性托管式集群,以处理持续性工作负载。
  • 优化的 BigQuery 连接:Spark-BigQuery 连接器直接以原生 Apache Arrow 格式使用 BigQuery 数据,从而绕过 JVM 面向行的转换。数据科学家还可以使用集成的笔记本界面,在同一个受治理的数据集上执行 PySpark DataFrame 代码和 SQL 查询,而无需切换环境。
  • 向量化原生 C++ 执行:在 Google Cloud 上运行 Spark 工作负载时,您可以为无服务器批处理或托管式集群启用 Lightning Engine。这个原生 C++ 查询执行引擎使用 Velox 和 Gluten 将物理查询计划直接编译为原生指令。它绕过了 JVM Volcano 迭代器模型和垃圾回收瓶颈,在不更改任何代码的情况下,可将 DataFrame 和 Spark SQL 工作负载的处理速度提升至原来的 4.9 倍。


更进一步

获享 $300 赠金以及 20 多种提供“始终免费”用量的产品,开始在 Google Cloud 上构建项目。