本文档介绍了如何设置和配置 Apache Spark 和 Apache Hive 以使用 Lakehouse 运行时目录。您将了解如何创建 Apache Hive 目录、配置 Spark 会话以连接到 metastore,以及运行工作负载以创建可以直接在 BigQuery 中查询的表。
准备工作
- 请参阅关于 Lakehouse 运行时目录中的 Hive 目录,了解 Spark 如何连接到 Lakehouse 运行时目录。
- 查看支持的存储格式和数据类型。
- 查看限制和注意事项。
- 登录您的 Google Cloud 账号。如果您是 Google Cloud新手,请 创建一个账号来评估我们的产品在实际场景中的表现。新客户还可获享 $300 赠金,用于运行、测试和部署工作负载。
-
Verify that billing is enabled for your Google Cloud project.
Enable the Lakehouse for Apache Iceberg, Dataproc API APIs.
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.-
Verify that billing is enabled for your Google Cloud project.
Enable the Lakehouse for Apache Iceberg, Dataproc API APIs.
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.
所需的角色
如需获得使用 Lakehouse 运行时目录所需的权限,请让您的管理员为您授予项目的以下 IAM 角色:
-
创建 Managed Service for Apache Spark 集群:
Dataproc 编辑器 (
roles/dataproc.editor) -
集群服务账号:
- Dataproc Worker (
roles/dataproc.worker) - Storage Object User (
roles/storage.objectUser) - BigLake Editor (
roles/biglake.editor) - Service Usage Consumer (
roles/serviceusage.serviceUsageConsumer)
- Dataproc Worker (
-
对所有 Lakehouse 运行时目录资源拥有写入权限:
BigLake 编辑者 (
roles/biglake.editor) -
对所有 Lakehouse 运行时目录资源的只读权限:BigLake Viewer (
roles/biglake.viewer)
如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限。
如需查看相关说明,请参阅授予单个角色。
常规工作流程
如需将 Lakehouse 运行时目录与 Spark 和 Hive 搭配使用,请按照以下常规工作流程操作:
- 为 Apache Iceberg Hive 目录创建湖仓一体。
- 使用您偏好的工具(例如 Managed Service for Apache Spark 或 BigQuery Studio)配置 Spark 会话。
- 在 Spark 会话中执行数据库和表操作。
- 向 Managed Service for Apache Spark 提交批量工作负载,并直接从 BigQuery 查询生成的表。
创建 Lakehouse Hive 目录
如需将 Lakehouse 运行时目录与 Spark 和 Hive 搭配使用,您必须先创建 Hive 目录。
Lakehouse Hive 目录是 Hive 数据库的集合。在运行 Spark 作业之前,请创建一个目录,以便将其注册到 Lakehouse Metastore。目录具有名称和 Cloud Storage 位置(Hive 数据所在的位置)。
控制台
在 Google Cloud 控制台中,打开 Lakehouse 页面。
点击创建目录。
选择 Lakehouse 运行时目录。
对于目录类型,选择 Hive Metastore。
在选择 Cloud Storage 存储桶字段中,输入要与目录搭配使用的 Cloud Storage 存储桶的名称。或者,点击浏览以选择现有存储桶或创建新存储桶。
对于 Catalog ID,为您的 Lakehouse Hive Catalog 命名。
对于主要位置,请指定与您的存储桶相同的区域。
点击创建。
gcloud
如需创建 Hive 目录,请运行以下命令:
gcloud beta biglake hive catalogs create LAKEHOUSE_CATALOG_ID \
--project=PROJECT_ID \
--location-uri="gs://GCS_WAREHOUSE_PATH" \
--primary-location=REGION \
--description="DESCRIPTION"
替换以下内容:
LAKEHOUSE_CATALOG_ID:Hive 目录名称。GCS_WAREHOUSE_PATH:用于存储 Hive 数据仓库的 Cloud Storage 路径。PROJECT_ID: Google Cloud项目 ID。REGION:元存储区的主区域。对于单区域存储桶,它应与存储桶区域一致。对于双区域或多区域存储桶,它应是组成区域之一,并且是元存储区主副本的预期位置。另一个区域将成为次要副本。DESCRIPTION:目录的说明。
curl
- 如需创建 Hive 目录,请运行以下命令:
curl -X POST -s -i -H "Authorization: Bearer $(gcloud auth print-access-token)" \
-d '{"locationUri": "gs://GCS_WAREHOUSE_PATH", "description": "DESCRIPTION"}' \
-H "Content-Type:application/json" \ "https://biglake.googleapis.com/hive/v1beta/projects/PROJECT_ID/catalogs?hiveCatalogId=LAKEHOUSE_CATALOG_ID&primary_location=REGION"
替换以下内容:
LAKEHOUSE_CATALOG_ID:Hive 目录名称。GCS_WAREHOUSE_PATH:用于存储 Hive 数据仓库的 Cloud Storage 路径。PROJECT_ID:您的 Google Cloud 项目 ID。REGION:元存储区的主要区域。对于单区域存储桶,它应与存储桶区域一致。对于双区域或多区域存储桶,它应是组成区域之一,也是元存储区的主要副本所在的区域。另一个区域将成为次要副本。DESCRIPTION:目录的说明。
配置和使用 Spark 及 Hive
如需使用 Lakehouse 运行时目录,您必须使用特定属性配置 Spark 会话。您可以在创建 Managed Service for Apache Spark 集群时设置这些属性,也可以在每次创建会话时指定这些属性。
这些属性包括客户端工厂、 Google Cloud项目 ID、默认目录和数据仓库目录等详细信息。会话建立后,您可以执行基本操作,例如列出现有数据库、创建新数据库、定义表和插入数据。
spark-sql
- 使用 SSH 连接到 Managed Service for Apache Spark 集群的主节点。
在命令行上运行
spark-sql并使用以下属性来启动交互式 Spark SQL 会话:spark-sql \ --conf spark.hive.metastore.client.factory.class=com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory \ --conf spark.hive.metastore.blms.project.id=PROJECT_ID \ --conf spark.hive.metastore.blms.catalog.default=LAKEHOUSE_CATALOG_ID \ --conf spark.hive.metastore.warehouse.dir=gs://GCS_WAREHOUSE_PATH会话开始后,Spark 会连接到 Lakehouse 运行时目录。
运行以下命令以创建和查询资源:
-- Show all the databases in the current project. SHOW DATABASES; -- Create a database. CREATE DATABASE spark_blms_database; -- Create a Parquet datasource table. CREATE TABLE spark_blms_database.parquet_quick_start (id INT, name STRING) USING PARQUET; -- Insert data into the table. INSERT INTO TABLE spark_blms_database.parquet_quick_start VALUES (1, 'my-first-user');替换以下内容:
PROJECT_ID:您的 Google Cloud 项目 ID。LAKEHOUSE_CATALOG_ID:Hive 目录名称。GCS_WAREHOUSE_PATH:用于存储 Hive 数据仓库的 Cloud Storage 路径。
Jupyter 笔记本
- 按照说明在 Managed Service for Apache Spark 集群上运行 Jupyter 笔记本。
- 在 Google Cloud 控制台的集群详情页面中,通过网页界面标签页访问 Jupyter 网页界面。
在新的笔记本中,创建 Spark 会话,然后运行以下查询:
from pyspark.sql import SparkSession # If a Spark session exists, stop it first by running spark.stop() spark = SparkSession.builder\ .master("local")\ .appName("Lakehouse runtime catalog tutorial")\ .config("spark.hive.metastore.client.factory.class", "com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory")\ .config("spark.hive.metastore.blms.project.id", "PROJECT_ID")\ .config("spark.hive.metastore.warehouse.dir", "gs://GCS_WAREHOUSE_PATH")\ .config("spark.hive.metastore.blms.catalog.default", "LAKEHOUSE_CATALOG_ID")\ .getOrCreate() # Show all the databases. df = spark.sql("SHOW DATABASES;") df.show() # Create a database. spark.sql("CREATE DATABASE jupyter_blms_db") # Create a Parquet datasource table. spark.sql("CREATE TABLE jupyter_blms_db.parquet_table(id INT, name STRING) USING PARQUET") # Insert data into the table. spark.sql("INSERT INTO TABLE jupyter_blms_db.parquet_table VALUES (1, 'my-first-user');") # Query from table. spark.sql("SELECT * FROM jupyter_blms_db.parquet_table;").show()替换以下内容:
PROJECT_ID:您的 Google Cloud 项目 ID。GCS_WAREHOUSE_PATH:用于存储 Hive 数据仓库的 Cloud Storage 路径。LAKEHOUSE_CATALOG_ID:Hive 目录名称。
BigQuery 笔记本
在 Google Cloud 控制台中,前往 BigQuery。
在探索器窗格中,点击 + 添加,然后点击 Python 笔记本。
在代码单元格中,配置 Managed Service for Apache Spark 会话和 Lakehouse 运行时目录属性:
from google.cloud.dataproc_spark_connect import DataprocSparkSession from google.cloud.dataproc_v1 import Session import os os.environ['DATAPROC_SPARK_CONNECT_DEFAULT_DATASOURCE'] = "" session = Session() session.environment_config.execution_config.ttl = {"seconds": 864000} session.runtime_config.version = "2.3" session.runtime_config.properties = { "spark.hive.metastore.blms.project.id": "PROJECT_ID", "spark.hive.metastore.blms.catalog.default": "LAKEHOUSE_CATALOG_ID", "spark.hive.metastore.warehouse.dir": "gs://GCS_WAREHOUSE_PATH", "spark.hive.metastore.client.factory.class": "com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory", "spark.sql.catalogImplementation": "hive" } spark = DataprocSparkSession.builder.dataprocSessionConfig(session).getOrCreate() print("Spark session created successfully")在另一个代码单元中,创建数据库和表:
# Create a database spark.sql("CREATE DATABASE bq_spark_blms_database;") # Create a parquet datasource table spark.sql("CREATE TABLE bq_spark_blms_database.parquet_quick_start (id INT, name STRING) USING PARQUET;") # Insert data into the table spark.sql("INSERT INTO TABLE bq_spark_blms_database.parquet_quick_start VALUES (1, 'my-first-user');") # Query from table spark.sql("select * from bq_spark_blms_database.parquet_quick_start;").show()替换以下内容:
PROJECT_ID:您的 Google Cloud 项目 ID。GCS_WAREHOUSE_PATH:用于存储 Hive 数据仓库的 Cloud Storage 路径。LAKEHOUSE_CATALOG_ID:Hive 目录名称。
Hive CLI
如需使用 Hive CLI,您必须将其配置为连接到 metastore。您可以通过以下两种方式之一来执行此操作:
- 选项 1:永久配置。更新主节点上的
/etc/hive/conf/hive-site.xml文件。 - 选项 2:临时配置。在启动 CLI 时提供配置标志。
如需配置 CLI 并运行查询会话,请按以下步骤操作:
- 使用 SSH 连接到 Managed Service for Apache Spark 集群的主节点。
根据您的配置方法,执行以下操作之一:
如需永久配置(选项 1):
在主节点上打开
/etc/hive/conf/hive-site.xml并添加以下属性:<property> <name>hive.metastore.blms.project.id</name> <value>PROJECT_ID</value> <description></description> </property> <property> <name>hive.metastore.client.factory.class</name> <value>com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory</value> <description></description> </property> <property> <name>hive.metastore.blms.catalog.default</name> <value>LAKEHOUSE_CATALOG_ID</value> <description></description> </property> <property> <name>hive.metastore.warehouse.dir</name> <value>gs://GCS_WAREHOUSE_PATH</value> <description></description> </property>启动 Hive CLI:
hive
在启动时临时配置(选项 2):使用配置标志启动 Hive CLI:
hive \ --hiveconf hive.metastore.blms.project.id="PROJECT_ID" \ --hiveconf hive.metastore.blms.catalog.default="LAKEHOUSE_CATALOG_ID" \ --hiveconf hive.metastore.client.factory.class=com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory \ --hiveconf hive.metastore.warehouse.dir="gs://GCS_WAREHOUSE_PATH"
然后,在 Hive CLI 会话中运行以下命令:
show databases; create database hive_query_test; create table hive_query_test.parquet_table (id INT, name STRING) stored as PARQUET; select * from hive_query_test.parquet_table;替换以下内容:
PROJECT_ID:您的 Google Cloud 项目 ID。LAKEHOUSE_CATALOG_ID:Hive 目录名称。GCS_WAREHOUSE_PATH:用于存储 Hive 数据仓库的 Cloud Storage 路径。
提交 Managed Service for Apache Spark 批量作业
您可以将 PySpark 批量工作负载提交到使用 Lakehouse 运行时目录的 Managed Service for Apache Spark。
此 PySpark 代码段会初始化一个配置为与 Lakehouse 运行时目录连接的 Spark 会话。它会设置 Lakehouse 客户端工厂、 Google Cloud 项目 ID、默认目录和仓库目录等基本属性。建立会话后,该代码演示了如何使用 Spark SQL 命令列出现有数据库、创建新数据库,以及在该数据库中定义 Parquet 格式的表。
创建一个包含 PySpark 作业的 Python 文件:
from pyspark.sql import SparkSession spark = ( SparkSession.builder.appName("Lakehouse runtime catalog tutorial") .config( "spark.hive.metastore.client.factory.class", "com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory") .config("spark.hive.metastore.blms.project.id", "PROJECT_ID") .config("spark.hive.metastore.blms.catalog.default", "LAKEHOUSE_CATALOG_ID") .config( "spark.hive.metastore.warehouse.dir", "gs://GCS_WAREHOUSE_PATH", ) .enableHiveSupport() .getOrCreate() ) # Show all the databases. spark.sql("SHOW DATABASES;").show() # Create a database. spark.sql("CREATE DATABASE dp_serverless_test") # Create a Parquet datasource table. spark.sql( "CREATE TABLE dp_serverless_test.parquet_table(id INT, name STRING) USING" " PARQUET" ) # Query from table. spark.sql("SELECT * from dp_serverless_test.parquet_table").show()替换以下内容:
PROJECT_ID:您的 Google Cloud 项目 ID。LAKEHOUSE_CATALOG_ID:Hive 目录名称。GCS_WAREHOUSE_PATH:用于存储 Hive 数据仓库的 Cloud Storage 路径。
提交批量作业:
gcloud dataproc batches submit pyspark PYTHON_SCRIPT_FILE \ --version=2.2 \ --project=PROJECT_ID \ --region=REGION \ --deps-bucket=gs://CLOUD_STORAGE_BUCKET替换以下内容:
PYTHON_SCRIPT_FILE:PySpark 应用文件的路径。可以是本地路径,也可以是 Cloud Storage 对象路径。PROJECT_ID:您的项目 ID。REGION:需要在其中运行批量作业的区域。CLOUD_STORAGE_BUCKET:用于暂存任何工作负载依赖项的 Cloud Storage 存储桶的名称。
查询 BigQuery 中的表
在 Lakehouse 运行时目录中通过 Spark 创建资源后,您可以在 BigQuery Studio 中查询这些资源。
在 Google Cloud 控制台中,前往 BigQuery。
在查询编辑器中,输入以下语句:
SELECT * FROM `PROJECT_ID.LAKEHOUSE_CATALOG_ID.DATABASE_NAME.TABLE_NAME`;替换以下内容:
PROJECT_ID:您的项目 ID。LAKEHOUSE_CATALOG_ID:您的 Hive 目录名称。DATABASE_NAME:您的数据库名称。TABLE_NAME:表格名称。