1. 简介 - Managed Service for Apache Spark
Managed Service for Apache Spark是一项具有高度可伸缩性的全代管式服务,用于运行 Apache Spark、Apache Flink、Presto 和众多其他开源工具和框架。使用 Managed Service for Apache Spark 可以大规模实现数据湖现代化改造、ETL / ELT 和安全数据科学。Managed Service for Apache Spark 还与多种 Google Cloud 服务完全集成,包括 BigQuery、Cloud Storage、Gemini Enterprise Agent Engine 和 Knowledge Catalog。
Managed Service for Apache Spark 提供两种部署模式:
- 借助托管式 Apache Spark Serverless,您无需配置基础设施和自动扩缩功能即可运行 PySpark 作业。托管式 Apache Spark 支持 PySpark 批量工作负载和会话 / 笔记本。
- 托管式 Apache Spark 集群:可让您为基于 YARN 的 Spark 工作负载以及 Flink 和 Presto 等开源工具管理 Hadoop YARN 集群。您可以根据需要纵向或横向扩缩云端集群,包括自动扩缩。
2. 在 Google Cloud VPC 上创建托管式 Apache Spark 集群
在此步骤中,您将使用 Google Cloud 控制台在 Google Cloud 上创建托管式 Apache Spark 集群。
首先,在控制台中启用 Managed Apache Spark 服务 API。启用后,在搜索栏中搜索“Managed Apache Spark”,然后点击创建集群。
选择 Compute Engine 上的集群,以使用 Google Compute Engine(GCE) 虚拟机作为运行 Managed Apache Spark 集群的基础架构。

您现在位于“集群创建”页面。

本页内容:
- 为集群提供一个唯一的名称。
- 选择特定区域。您也可以选择地区,不过,Managed Apache Spark 能够自动为您选择地区。在此 Codelab 中,请选择“us-central1”和“us-central1-c”。
- 选择“标准”集群类型。这样可确保有一个主节点。
- 在配置节点标签页中,确认创建的工作器数量为 2。
- 在自定义集群部分中,选中启用组件网关旁边的复选框。这样,您就可以访问集群上的网页界面,包括 Spark 界面、Yarn 节点管理器和 Jupyter 笔记本。
- 在可选组件中,选择 Jupyter 笔记本。这会为集群配置 Jupyter 笔记本服务器。
- 将所有其他设置保留原样,然后点击创建集群。
这会启动一个托管式 Apache Spark 集群。
3. 启动集群并通过 SSH 登录集群
当集群状态变为正在运行后,在 Managed Apache Spark 控制台中点击集群名称。

点击虚拟机实例标签页,查看集群的主节点和两个工作器节点。

点击主节点旁边的 SSH 以登录主节点。

运行 hdfs 命令以查看目录结构。
hadoop_commands_example
sudo hadoop fs -ls /
sudo hadoop version
sudo hadoop fs -mkdir /test51
sudo hadoop fs -ls /
4. 网页界面和组件网关
在 Managed Apache Spark 集群控制台中,点击集群的名称,然后点击 WEB 界面标签页。

这会显示可用的 Web 界面,包括 Jupyter。点击 Jupyter 以打开 Jupyter 笔记本。您可以使用此笔记本在 GCS 中存储的 PySpark 中创建笔记本。将笔记本存储在 Google Cloud Storage 中,然后打开 PySpark 笔记本以在此 Codelab 中使用。
5. 监控和观察 Spark 作业
在托管式 Apache Spark 集群正常运行后,创建 PySpark 批量作业,并将该作业提交到托管式 Apache Spark 集群。
创建一个 Google Cloud Storage (GCS) 存储分区,用于存储 PySpark 脚本。请务必在托管式 Apache Spark 集群所在的区域中创建该存储分区。

现在,GCS 存储分区已创建完毕,请将以下文件复制到此存储分区中。
https://raw.githubusercontent.com/diptimanr/spark-on-gce/main/test-spark-1.py
此脚本会创建一个示例 Spark DataFrame,并将其写入为 Hive 表。
hive_job.py
from pyspark.sql import SparkSession
from datetime import datetime, date
from pyspark.sql import Row
spark = SparkSession.builder.master("local").enableHiveSupport().getOrCreate()
df = spark.createDataFrame([ (1, 2., 'string1', date(2000, 1, 1), datetime(2000, 1, 1, 12, 0)),
(2, 3., 'string2', date(2000, 2, 1), datetime(2000, 1, 2, 12, 0)), (3, 4., 'string3', date(2000, 3, 1), datetime(2000, 1, 3, 12, 0))
], schema='a long, b double, c string, d date, e timestamp')
print("..... Writing data .....")
df.write.mode("overwrite").saveAsTable("test_table_1")
print("..... Complete .....")
在 Managed Apache Spark 中将此脚本作为 Spark 批量作业提交。点击左侧导航菜单中的作业,然后点击提交作业

提供任务 ID 和区域。选择您的集群,并提供您复制的 Spark 脚本的 GCS 位置。此作业将作为 Spark 批量作业在 Managed Apache Spark 上运行。
在属性下,添加键 spark.submit.deployMode 和值 client,以确保驱动程序在 Managed Apache Spark 主节点中运行,而不是在工作器节点中运行。点击提交,将批量作业提交到 Managed Apache Spark。

Spark 脚本将创建一个 DataFrame 并写入 Hive 表 test_table_1。
作业成功运行后,您可以在监控标签页下看到控制台打印语句。

现在,Hive 表已创建完毕,请提交另一个 Hive 查询作业,以选择该表的内容并在控制台上显示。
创建另一个具有以下属性的作业:

请注意,作业类型设置为 Hive,查询源类型为查询文本,这意味着我们将在查询文本文本框中编写整个 HiveQL 语句。
提交作业,其余参数保留为默认值。

请注意,HiveQL 如何选择所有记录并显示在控制台上。
6. 自动扩缩
自动扩缩是指估算工作负载的“适当”集群工作器节点数量的任务。
Managed Apache Spark AutoscalingPolicies API 提供自动管理集群资源的机制,还启用了集群工作器虚拟机的自动扩缩功能。自动扩缩政策是可重复使用的配置,描述了应如何扩缩使用该自动扩缩政策的集群工作器。该政策定义了扩缩边界、频率和积极性,以在整个集群生命周期内提供对集群资源的精细控制。
托管式 Apache Spark 自动扩缩政策使用 YAML 文件编写,这些 YAML 文件在创建集群时通过 CLI 命令传入,或者在通过 Cloud 控制台创建集群时从 GCS 存储分区中选择。
以下是托管式 Apache Spark 自动扩缩政策的示例:
policy.yaml
workerConfig:
minInstances: 10
maxInstances: 10
secondaryWorkerConfig:
maxInstances: 50
basicAlgorithm:
cooldownPeriod: 4m
yarnConfig:
scaleUpFactor: 0.05
scaleDownFactor: 1.0
gracefulDecommissionTimeout: 1h
7. 配置托管式 Apache Spark 可选组件
这会启动一个托管式 Apache Spark 集群。
创建 Managed Apache Spark 集群时,标准 Apache Hadoop 生态系统组件会自动安装在集群中(请参阅 Managed Apache Spark 版本列表)。您可以在创建集群时在集群上安装称为可选组件的其他组件。

在通过控制台创建托管式 Apache Spark 集群时,我们已启用可选组件,并选择 Jupyter Notebook 作为可选组件。
8. 清理资源
如需清理集群,请在 Managed Apache Spark 控制台中选择集群后,点击停止。集群停止后,点击删除以删除集群。
删除 Managed Apache Spark 集群后,删除复制了代码的 GCS 存储分区。
如需清理资源并避免产生任何不必要的结算费用,您需要先停止再删除受管 Apache Spark 集群。
在停止并删除集群之前,请确保已将写入 HDFS 存储空间的所有数据复制到 GCS 以进行持久存储。
如需停止集群,请点击停止。

集群停止后,点击删除以删除集群。
在确认对话框中,点击删除以删除集群。
