使用 Lightning Engine
Lightning Engine 是新一代 Apache Spark 性能,引入了独有的增强功能,旨在显著提升性能、成本效益和运营稳定性。
优势
Lightning Engine 的优势包括:
加速数据操作:通过优化与云存储的互动(包括元数据处理、写入工作负载和矢量化 I/O),显著提升性能并节省成本。
智能查询执行:利用高级优化器增强功能,动态减少扫描的数据量、优化数据处理,并生成更高效的执行计划,从而实现更快、更具成本效益的查询。
简化 AI 和机器学习工作负载:缩短基于 GPU 的工作负载的集群启动时间,并使用针对 AI 优化的映像简化安全环境中的部署。
虽然 Lightning Engine 可大幅提升性能,但具体影响因工作负载而异。它最适合利用 Spark Dataframe API、Spark Dataset API 和 Spark SQL 查询的计算密集型任务,而不是 I/O 密集型操作。
与标准引擎的比较
Lightning Engine 是标准引擎的替代方案,用于在 Managed Service for Apache Spark 集群上执行 Spark 作业。下表比较了 Lightning Engine 与标准引擎的激活属性、工作负载适用性和主要优势。
| 功能 | 标准引擎 | Lightning Engine |
|---|---|---|
| CLI 标志 | --engine=default 或取消设置标记 |
--engine=lightning |
| 适用场景 | 一般用途作业、开发和测试 | 需要大幅加速的企业级工作负载 |
| 主要优势 | 基准性能 | 优化的云存储互动,智能查询执行 |
要求
以下要求适用于 Lightning Engine 功能:
- 映像版本:Lightning Engine 必须与 Managed Service for Apache Spark 映像版本
2.3.3及更高版本的2.3子次要映像版本搭配使用。 Managed Service for Apache Spark 映像版本3.0不支持 Lightning Engine。 - 支持的作业:支持 Spark、PySpark、SparkSQL 和 SparkR。标准引擎在提交给 Lightning Engine 集群的其他作业类型上运行。
原生查询执行
原生查询执行 (NQE) 是 Lightning Engine 的可选组件,可为特定作业提供更深层次的加速。它是一款基于 Apache Gluten 和 Velox 的原生引擎,针对 Google 硬件进行了优化,可通过在 JVM 外部运行部分 Spark 查询来提升性能。
- 建议将 NQE 用于以下方面:
- 利用 Spark Dataframe API 和 Spark Dataset API 的计算密集型任务,以及从 Parquet、ORC、Apache Iceberg 和 Delta Lake 文件和表中读取数据的 Spark SQL 查询。输出文件格式不会影响其性能。
- 不建议将 NQE 用于以下方面:
- 严重依赖弹性分布式数据集 (RDD)、用户定义的函数 (UDF)、大多数 Spark 机器学习 (ML) 库以及因存储空间访问而延迟的 I/O 绑定操作的作业。
ARM 上的原生查询执行
Managed Service for Apache Spark 支持在 ARM 架构上的 Lightning Engine 中执行原生查询,专门针对 C4A ARM 虚拟机上的 Google Axion 处理器进行了优化(请参阅 Managed Service for Apache Spark 支持的 ARM 机器类型),并使用 2.3-ubuntu22-arm 映像。Google Axion C4A 实例是基于 ARM 的虚拟机,基于 ARM Neoverse V2 核心构建,可显著提升性价比。
在 Google Axion ARM 实例上运行 NQE 可将硬件级性价比和高单核内存带宽与软件级原生向量化查询加速相结合,从而显著缩短查询运行时并降低总拥有成本 (TCO)。
ARM 的主要优势包括:
- 速度比开源 Spark 快 2.5 倍到 4 倍:与开源 Apache Spark 相比,在分析型工作负载方面,性能提升了 2.5 倍到 4 倍。
- 费用优化:最大限度地发挥 Google Axion C4A 实例的性价比优势(与同类 x86 实例相比,性价比提升高达 65%)。
- 完全对等的运算符和格式:与 x86 NQE 运算符完全对等,可为 Cloud Storage Parquet、Apache Iceberg 和 Delta Lake 表格式提供原生加速。
- 无需更改任何代码:使用标准配置标志在 ARM 集群上无缝激活,无需修改 Spark 应用代码。
要求
以下要求适用于原生查询执行功能:
执行引擎:NQE 仅适用于在创建集群时启用了 Lightning 引擎的集群。
支持的作业:支持 Spark、PySpark、SparkSQL 和 SparkR。标准引擎将在提交到 Lightning Engine 集群的其他作业类型上运行(不使用 NQE)。
无 GPU 和加速器:在 GPU 加速器上提交的启用 NQE 的作业将失败(但可以在没有 NQE 的情况下受益于 Lightning Engine)。
数据类型:不支持以下数据类型的输入:
- 字节:ORC 和 Parquet
- 结构体、数组、Map:Parquet
架构和操作系统要求
NQE 要求因处理器架构而异:
| 架构 | 支持的机器类型 | 支持的操作系统和映像 |
|---|---|---|
| x86 | Intel 和 AMD 机器家族 | Debian-12 和 Ubuntu-22(映像版本 2.3.3 及更高版本的 2.3 版本) |
| ARM | C4A 机器系列(Google Axion ARM 虚拟机) | 2.3-ubuntu22-arm 映像(版本 2.3.3 及更高版本的 2.3 发布版本) |
价格
如需了解价格信息,请参阅 Managed Service for Apache Spark 价格。
创建 Lightning Engine 集群
本部分介绍如何创建 Managed Service for Apache Spark 集群,以在提交到该集群的 Spark 作业上启用 Lightning Engine。
您还可以在创建集群时为集群启用原生查询执行 (NQE),也可以稍后为提交到集群的特定 Spark 作业启用 NQE。
准备工作
-
In the Google Cloud console, on the project selector page, select or create a Google Cloud project.
Roles required to select or create a project
- Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
-
Create a project: To create a project, you need the Project Creator role
(
roles/resourcemanager.projectCreator), which contains theresourcemanager.projects.createpermission. Learn how to grant roles.
-
Verify that you have the permissions required to complete this guide.
-
Verify that billing is enabled for your Google Cloud project.
Enable the Dataproc API.
Roles required to enable APIs
To enable APIs, you need the
serviceusage.services.enablepermission. If you created the project, then you likely already have this permission through the Owner role (roles/owner). Otherwise, you can get this permission through the Service Usage Admin role (roles/serviceusage.serviceUsageAdmin). Learn how to grant roles.-
安装 Google Cloud CLI。
-
配置 gcloud CLI 以使用您的联合身份。
如需了解详情,请参阅使用联合身份登录 gcloud CLI。
-
如需初始化 gcloud CLI,请运行以下命令:
gcloud init
所需的角色
您需要拥有特定的 Identity and Access Management (IAM) 角色才能创建 Managed Service for Apache Spark 集群并向该集群提交作业。根据组织政策,云项目所有者或服务管理员可能已向您或服务账号授予这些角色。如需查看角色授予情况,请参阅您是否需要授予角色?。
如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限。
用户角色
如需获得创建 Managed Service for Apache Spark 集群所需的权限,请让管理员向您授予以下 IAM 角色:
- 全部:
- 针对项目的 Dataproc Editor (
roles/dataproc.editor) - Compute Engine 默认服务账号的 Service Account User (
roles/iam.serviceAccountUser)
- 针对项目的 Dataproc Editor (
服务账号角色
如需确保 Compute Engine 默认服务账号具有创建 Managed Service for Apache Spark 集群所需的权限,请让您的管理员向 Compute Engine 默认服务账号授予项目的 Dataproc Worker (roles/dataproc.worker) IAM 角色。
创建集群
以下示例展示了如何使用 Google Cloud 控制台、Google Cloud CLI、Dataproc API、适用于 Python 的 Cloud 客户端库或 Terraform 创建 Lightning Engine 集群。您还可以使用 Go、Java 和 Node.js Cloud 客户端库创建启用 Lightning Engine 的集群。
控制台
- 打开创建集群页面。
- 点击其他配置以展开该部分。
- 修改自定义及其他。
- 在随即打开的面板中,确认已选中启用 Lightning Engine 复选框。
- 可选:如需为 Spark 作业默认启用原生执行运行时,请选中启用原生执行复选框。
- 点击保存。
- 根据需要配置其他集群设置。
- 点击创建集群。
创建 ARM 集群:
如需创建具有 Google Axion 和原生查询执行功能的 ARM 集群,请执行以下操作:
- 在工作器配置中,选择 C4A 机器系列,然后选择您的 ARM 机器类型。
- 在图片中,选择
2.3-ubuntu22-arm图片。 - 在其他配置 > 自定义和其他中,确认已选择启用 Lightning 引擎和启用原生执行。
- 点击创建集群。
gcloud
如需创建启用了 Lightning Engine 的集群,请运行
gcloud dataproc clusters create命令并添加--engine=lightning标志。如需了解详情,请参阅使用 gcloud CLI 创建集群。gcloud dataproc clusters create CLUSTER_NAME \ --region=REGION \ --engine=lightning \ --image-version=2.3可选:如需默认针对 Spark 作业启用原生执行运行时,请添加
spark:spark.dataproc.lightningEngine.runtime=native属性。gcloud dataproc clusters create CLUSTER_NAME \ --region=REGION \ --engine=lightning \ --image-version=2.3 \ --properties='spark:spark.dataproc.lightningEngine.runtime=native'创建启用 NQE 的 ARM 集群:如需创建启用 NQE 的 Managed Service for Apache Spark 集群,并使用 Google Axion C4A ARM 节点,请执行以下操作:
gcloud dataproc clusters create CLUSTER_NAME \ --region=REGION \ --image-version=2.3-ubuntu22-arm \ --master-machine-type=c4a-standard-16 \ --worker-machine-type=c4a-standard-16 \ --engine=lightning \ --properties='spark:spark.dataproc.lightningEngine.runtime=native'
API
如需创建启用了 Lightning Engine 的集群,请发送 clusters.create 请求。如需了解详情,请参阅使用 REST API 创建集群。
在请求正文中,将
engine字段设置为LIGHTNING。{ "projectId": "PROJECT_ID", "clusterName": "CLUSTER_NAME", "config": { "engine": "LIGHTNING", "gceClusterConfig": {}, "softwareConfig": { "imageVersion": "2.3" } } }可选:如需默认针对所有作业启用原生执行运行时,请添加
spark:spark.dataproc.lightningEngine.runtime属性。{ "projectId": "PROJECT_ID", "clusterName": "CLUSTER_NAME", "config": { "engine": "LIGHTNING", "gceClusterConfig": {}, "softwareConfig": { "imageVersion": "2.3", "properties": { "spark:spark.dataproc.lightningEngine.runtime": "native" } } } }
Python
如需创建启用了 Lightning Engine 的集群,请使用
create_cluster方法,并将集群配置中的engine字段设置为LIGHTNING。如需了解详情,请参阅使用 Python 创建集群。from google.cloud import dataproc_v1 def create_lightning_cluster(project_id, region, cluster_name): client_options = {"api_endpoint": f"{region}-dataproc.googleapis.com:443"} cluster_client = dataproc_v1.ClusterControllerClient(client_options=client_options) cluster = { "project_id": project_id, "cluster_name": cluster_name, "config": { "engine": "LIGHTNING", "software_config": { "image_version": "2.3-debian12", }, } } operation = cluster_client.create_cluster( project_id=project_id, region=region, cluster=cluster ) result = operation.result() print(f"Cluster created successfully: {result.cluster_name}")可选:如需默认针对 Spark 作业启用原生执行运行时,请添加
spark:spark.dataproc.lightningEngine.runtime属性。from google.cloud import dataproc_v1 def create_lightning_native_cluster(project_id, region, cluster_name): client_options = {"api_endpoint": f"{region}-dataproc.googleapis.com:443"} cluster_client = dataproc_v1.ClusterControllerClient(client_options=client_options) cluster = { "project_id": project_id, "cluster_name": cluster_name, "config": { "engine": "LIGHTNING", "software_config": { "image_version": "2.3-debian12", "properties": { "spark:spark.dataproc.lightningEngine.runtime": "native" } } } } operation = cluster_client.create_cluster( project_id=project_id, region=region, cluster=cluster ) result = operation.result() print(f"Cluster created successfully: {result.cluster_name}")
Terraform
- 在
google_dataproc_cluster资源配置中,将engine实参设置为LIGHTNING。 - 如需了解详情和高级选项,请参阅
google_dataproc_cluster资源的官方 Terraform 文档。
验证集群引擎
控制台
- 在 Google Cloud 控制台中,前往集群详情页面。
- 验证
Lightning Engine值是否列在 Engine 字段中。 - 如果您启用了原生查询执行,请验证 原生执行字段中是否列出了
native。
gcloud
如需验证引擎和 NQE(如果已启用),请运行
gcloud dataproc clusters describe命令:gcloud dataproc clusters describe CLUSTER_NAME --project=PROJECT_ID --region=REGION检查输出中是否包含
engine和lightningEngine.runtime属性:clusterName: lightning-engine-cluster engine: lightningEngine lightningEngine.runtime: native
使用 Lightning Engine 提交作业
如果您在创建集群时启用了 Lightning Engine,那么当您向集群提交 Spark 作业时,系统默认会在该作业上启用 Lightning Engine。
为作业启用原生查询执行
如果您在创建 Lightning Engine 集群时启用了原生查询执行 (NQE),则所有 Spark 作业都会在启用 NQE 的情况下运行,除非您针对特定作业停用 NQE。
如果您在创建 Lightning Engine 集群时未启用 NQE,则可以在提交作业时为作业启用 NQE,如以下示例所示。
gcloud
如需在提交 Spark 作业时启用原生查询执行,请添加 spark.dataproc.lightningEngine.runtime=native 属性:
gcloud dataproc jobs submit spark \
--cluster=CLUSTER_NAME \
--region=REGION \
--properties=spark.dataproc.lightningEngine.runtime=native \
-- ...
API
如需在提交 Spark 作业时启用原生查询执行,请在请求中添加 spark.dataproc.lightningEngine.runtime 属性:
{
"job":{
"placement":{
"clusterName": ...
},
"sparkJob":{
"mainClass": ...,
"properties":{
"spark.dataproc.lightningEngine.runtime":"native"
}
}
}
}
为作业停用原生查询执行
如果您在创建 Lightning Engine 集群时启用了原生查询执行 (NQE),则所有 Spark 作业都将在启用 NQE 的情况下运行,除非您针对特定作业停用 NQE。
您可以在提交作业时针对特定 Spark 作业停用 NQE,如以下示例所示。
gcloud
如需在向 Lightning Engine 集群提交 Spark 作业时停用原生查询执行,请添加 spark.dataproc.lightningEngine.runtime=default 属性:
gcloud dataproc jobs submit spark \
--cluster=CLUSTER_NAME \
--region=REGION \
--properties=spark.dataproc.lightningEngine.runtime=default \
-- ...
API
如需在向 Lightning Engine 集群提交 Spark 作业时停用原生查询执行,请添加 spark.dataproc.lightningEngine.runtime=default 属性:
{
"job":{
"placement":{
"clusterName": ...
},
"sparkJob":{
"mainClass": ...,
"properties":{
"spark.dataproc.lightningEngine.runtime":"default"
}
}
}
}
验证作业的原生查询执行
向 Lightning Engine 集群提交作业后,您可以验证该作业是否已启用原生查询执行。
控制台
- 在 Google Cloud 控制台中,前往作业页面。
- 点击作业 ID 以打开作业详情页面。
- 验证 Native Execution 字段中是否列出了
native。
gcloud
运行
gcloud dataproc jobs describe命令:gcloud dataproc jobs describe JOB_ID --project=PROJECT_ID --region=REGION在属性部分中,检查输出中的
lightningEngine.runtime:lightningEngine.runtime: native
配置参数
下表总结了 Lightning Engine 和原生查询执行的主要配置参数。
| 参数名称 | 说明 | 适用的引擎 | 默认值 | 默认值(Lightning Engine) | 用户可覆盖(作业级) | 范围 |
|---|---|---|---|---|---|---|
--engine |
集群级设置,用于在创建集群期间选择引擎。 | 集群范围 | default |
lightning |
否 | 集群 |
spark:spark.dataproc.lightningEngine.runtime |
在创建集群期间选择 Lightning 引擎运行时所需的集群级设置。 | 仅限 Lightning | default |
default |
否 | 集群 |
spark.dataproc.lightningEngine.runtime |
启用或停用 Lightning Engine 中的原生查询执行 (NQE)。 | 仅限 Lightning | default |
default |
可以。可以设置为 native 或 default。 |
作业 |
限制
在以下场景中启用原生查询执行可能会导致异常、Spark 不兼容或工作负载回退到默认的 Spark 引擎。
后备广告
在以下情况下,原生查询执行可能会导致工作负载回退到 Spark 执行引擎:
- ANSI:如果启用 ANSI 模式,执行会回退到 Spark。
- 区分大小写模式:原生查询执行仅支持 Spark 默认的不区分大小写模式。如果启用了区分大小写模式,可能会出现错误的结果。
- 分区表扫描:仅当路径包含分区信息时,原生查询执行才支持分区表扫描。否则,工作负载会回退到 Spark 执行引擎。
不兼容的行为
在以下情况下使用原生查询执行时,可能会出现不兼容的行为或错误的结果:
- JSON 函数:原生查询执行支持用英文双引号而非英文单引号括起来的字符串。使用单引号时会出现错误结果。在路径中使用
*和get_json_object函数会返回NULL。 - Parquet 读取配置:
- 即使将
spark.files.ignoreCorruptFiles设置为true,原生查询执行也会将spark.files.ignoreCorruptFiles视为设置为默认值false。 - 原生查询执行会忽略
spark.sql.parquet.datetimeRebaseModeInRead,仅返回 Parquet 文件内容。不考虑旧版混合日历与 Proleptic 格里高利历之间的差异。Spark 结果可能会有所不同。
- 即使将
- NaN:不支持。例如,在数值比较中使用
NaN时,可能会出现意外结果。 - Spark 列式读取:由于 Spark 列式向量与原生查询执行不兼容,可能会发生严重错误。
- 溢出:当您将 shuffle 分区设置为较大数量时,溢出到磁盘功能可能会触发
OutOfMemoryException。如果发生这种情况,减少分区数量可以消除此异常。