Apache Spark 是一个分布式处理引擎,旨在大规模处理结构化、半结构化和非结构化数据。为了与引擎交互并定义分布式数据操作,开发者会使用 Spark API。
现代 Spark 开发的核心是 Spark DataFrame API。DataFrame 以具名列的形式整理和处理分布式数据,结构类似于关系型数据库中的表或电子表格。它提供了一个结构化的声明式接口,用于分析海量数据集,同时在后台自动优化物理执行。
如需了解 Spark API 如何在分布式网络中执行代码,必须先了解其运行时架构的基本组件:
在设计分布式数据流水线时,开发者需要在两种主要编程接口之间做出选择:
对于几乎所有数据工程和数据科学应用场景,DataFrame API 都是首选。RDD 仅适用于必须使用自定义 Java 对象处理原始非结构化数据(如二进制文件或媒体流)的场景。DataFrame 需要的代码更少,自动执行速度更快,并利用堆外二进制内存管理来绕过 JVM 垃圾回收瓶颈。
Spark DataFrame API 简化了处理、清理和准备大量数据的计算密集型任务。
Spark DataFrame API 是 Spark 高级处理库的编程基础:
声明式优化
DataFrame 使用内置的查询优化器,称为 Catalyst 优化器。当您编写 DataFrame 代码时,优化器会自动重写物理执行计划,应用过滤条件下推、投影剪枝和联接重新排序,从而在无需手动调优的情况下尽可能快地运行。
语言灵活性
无论您的团队使用 Python (PySpark)、Scala、Java 还是 R 编写代码,DataFrame API 都能确保相同的性能。底层执行计划在同一个独立于 JVM 的执行层中进行编译和优化。
高效的内存管理
DataFrame 利用 Tungsten 执行引擎,以高度压缩的堆外二进制格式存储数据。这消除了 JVM 对象创建开销,并防止垃圾回收暂停冻结执行器线程。
传统上,运行开源 Spark 需要手动配置集群、管理软件版本以及进行复杂的网络对等互连。Google Cloud 通过提供 Managed Service for Apache Spark 简化了这一过程,将分布式执行转换为全托管式、企业级数据平台。
数据团队使用以下功能在 Google Cloud 上执行 Spark DataFrame 工作负载: