本教程演示了如何从 Cloud Storage 上的 Parquet 文件中读取纽约市 (NYC) 出租车行程记录,使用 Managed Service for Apache Spark 无服务器运行时 3.0 上的 Apache Spark 对其进行转换,然后使用 BigLake REST 目录将汇总结果写入 Apache Iceberg 表中。
在此教程中,您需要完成以下任务:
- 暂存纽约市出租车 Parquet 数据。
- 编写 PySpark 作业脚本。
- 提交无服务器批量作业。
- 验证并查询 Iceberg 表。
准备工作
设置项目并执行其他启动任务。
设置 Google Cloud 项目
根据需要设置项目,以启用 API、授予 Identity and Access Management (IAM) 角色、验证应用默认凭据并创建 Cloud Storage 存储桶。
启用 API
使用 Google Cloud 控制台启用所需的 API。
-
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 billing is enabled for your Google Cloud project.
Enable the Dataproc, Cloud Storage, BigQuery, BigLake, Cloud Logging, Compute Engine, Cloud Resource Manager, and Dataproc Resource Manager 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.-
为您的用户账号创建本地身份验证凭证:
gcloud auth application-default login
如果系统返回身份验证错误,并且您使用的是外部身份提供方 (IdP),请确认您已 使用联合身份登录 gcloud CLI。
-
创建 Cloud Storage 存储桶:
将gcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAME替换为符合存储桶命名要求的存储桶名称。
根据需要授予 IAM 角色
您需要具备某些 IAM 角色才能运行本页中的示例。 这些角色可能已获授予,具体取决于组织政策。 如需查看角色授予情况,请参阅是否需要授予角色?。
如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限。
用户角色
默认情况下,Managed Service for Apache Spark 无服务器运行时 3.0 在最终用户凭证 (EUC) 下运行。不需要服务账号角色。如需了解详情,请参阅角色和无服务器 IAM 角色。
如需获得提交无服务器批处理工作负载所需的权限,请让您的管理员为您授予以下 IAM 角色:
-
在运行时 3.x(默认 EUC)上运行工作负载:
- 针对项目的 Dataproc Serverless Editor (
roles/dataproc.serverlessEditor) - 针对项目的 BigLake Admin (
roles/biglake.admin) 角色 - 针对项目的 BigQuery Admin (
roles/bigquery.admin) 角色 - 项目的 Storage Object Admin (
roles/storage.objectAdmin)
- 针对项目的 Dataproc Serverless Editor (
您的管理员可以运行以下 bash 脚本,以向您的用户账号授予角色。
```bash
for ROLE in \
roles/dataproc.serverlessEditor \
roles/biglake.admin \
roles/bigquery.admin \
roles/storage.objectAdmin
do
gcloud projects add-iam-policy-binding "${PROJECT_ID}" \
--member="user:${USER_ACCOUNT}" \
--role="${ROLE}"
done
```
配置环境变量
运行以下 bash 脚本,以设置本教程中使用的 shell 环境变量。
# 1. Active project ID, user account, and project number
export PROJECT_ID="$(gcloud config get-value project)"
export USER_ACCOUNT="$(gcloud config get-value account)"
export PROJECT_NUMBER="$(gcloud projects describe "${PROJECT_ID}" --format="value(projectNumber)")"
# 2. Regional deployment and networking settings
export REGION="us-central1"
export SUBNET_NAME="default"
export SUBNET_RANGE="10.128.0.0/20"
# 3. Storage and BigLake Iceberg catalog resources
export BUCKET_NAME="BUCKET_NAME"
export OUTPUT_CATALOG_NAME="lakehouse"
export OUTPUT_DATASET_NAME="nyc_taxi"
export OUTPUT_TABLE_NAME="yellow_trips_analyzed"
export INPUT_PARQUET_PATH="gs://${BUCKET_NAME}/raw/nyc_taxi/yellow_tripdata_2024-01.parquet"
替换以下内容:
BUCKET_NAME:您在设置 Google Cloud 项目中创建的 Cloud Storage 存储桶的名称。
第 1 步:创建 BigLake Iceberg 目录
借助 BigLake Iceberg REST Catalog,Spark 工作负载和 BigQuery 可以发现、读取和写入 Iceberg 表。
创建 BigLake Iceberg 目录。
创建目录,并将其默认位置指向 Cloud Storage 中的数据仓库路径。
gcloud biglake iceberg catalogs create "${OUTPUT_CATALOG_NAME}" \ --project="${PROJECT_ID}" \ --catalog-type=biglake \ --default-location="gs://${BUCKET_NAME}/warehouse" \ --credential-mode=end-user验证目录健康状况。
gcloud biglake iceberg catalogs describe "${OUTPUT_CATALOG_NAME}" --project="${PROJECT_ID}"
第 2 步:暂存纽约市出租车 Parquet 数据
纽约市出租车和豪华轿车委员会 (TLC) 每月都会以 Parquet 格式发布行程记录。下载一个月的数据并将其复制到您的 Cloud Storage 存储桶。
curl -L -o yellow_tripdata_2024-01.parquet \
https://d37ci6vzurychx.cloudfront.net/trip-data/yellow_tripdata_2024-01.parquet
gcloud storage cp yellow_tripdata_2024-01.parquet "${INPUT_PARQUET_PATH}"
gcloud storage ls "${INPUT_PARQUET_PATH}"
第 3 步:编写 PySpark 作业脚本
创建 taxi_to_iceberg.py。此脚本会读取原始 Parquet 记录,过滤无效行程,计算每日乘客数、票价、小费和收入汇总,并使用 BigLake REST 目录写入 Iceberg 表。
"""Read NYC taxi Parquet from Cloud Storage and write aggregates to an Iceberg table."""
import sys
from pyspark.sql import SparkSession
from pyspark.sql.functions import (
avg,
col,
count,
round as spark_round,
sum as spark_sum,
to_date,
when,
)
def main():
if len(sys.argv) < 5:
print(
"Usage: taxi_to_iceberg.py <input_parquet_path> <catalog_name>"
" <dataset_name> <table_name>"
)
sys.exit(1)
input_path, catalog_name, dataset_name, table_name = sys.argv[1:5]
full_table = f"`{catalog_name}`.{dataset_name}.{table_name}"
# Iceberg extensions are enabled here; catalog properties are passed at
# submit time via --flags-file (see Step 4).
spark = (
SparkSession.builder.appName("NYC Taxi Parquet to Iceberg")
.config(
"spark.sql.extensions",
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions",
)
.getOrCreate()
)
print(f"Reading raw Parquet data from: {input_path}")
raw_df = spark.read.parquet(input_path)
raw_df.printSchema()
cleaned_df = (
raw_df.filter(
(col("trip_distance") > 0)
& (col("fare_amount") > 0)
& (col("passenger_count") > 0)
)
.withColumn("trip_date", to_date(col("tpep_pickup_datetime")))
.withColumn(
"tip_pct",
when(
col("fare_amount") > 0,
spark_round((col("tip_amount") / col("fare_amount")) * 100, 2),
).otherwise(0.0),
)
)
aggregated_df = (
cleaned_df.groupBy("trip_date", "payment_type")
.agg(
count("*").alias("total_trips"),
spark_sum("passenger_count").alias("total_passengers"),
spark_round(avg("trip_distance"), 2).alias("avg_distance"),
spark_round(avg("fare_amount"), 2).alias("avg_fare"),
spark_round(avg("tip_pct"), 2).alias("avg_tip_percentage"),
spark_round(spark_sum("total_amount"), 2).alias("total_revenue"),
)
.orderBy("trip_date", "payment_type")
)
aggregated_df.show(10, truncate=False)
# Create the namespace, then write with DataFrameWriterV2 (writeTo).
spark.sql(f"CREATE NAMESPACE IF NOT EXISTS `{catalog_name}`.{dataset_name}")
spark.catalog.setCurrentCatalog(catalog_name)
aggregated_df.writeTo(full_table).using("iceberg").createOrReplace()
print(f"Wrote Iceberg table: {full_table}")
spark.stop()
if __name__ == "__main__":
main()
重要提示:在写入 BigLake Iceberg 目录时,请务必使用 writeTo (DataFrameWriterV2) API,先设置当前目录,确保命名空间存在,并在 SQL 和 writeTo 调用中用反引号将目录名称括起来。
第 4 步:提交无服务器批量作业
在配置文件中定义目录属性,并将批量作业提交到 Managed Service for Apache Spark 无服务器。
在 YAML 标志文件中定义 Iceberg catalog 属性。
创建包含 BigLake REST 目录属性的
iceberg-flags.yaml。Google Cloud CLI 会使用--flags-file实参读取此文件。cat <<EOF > iceberg-flags.yaml --properties: spark.sql.catalog.${OUTPUT_CATALOG_NAME}: org.apache.iceberg.spark.SparkCatalog spark.sql.catalog.${OUTPUT_CATALOG_NAME}.type: rest spark.sql.catalog.${OUTPUT_CATALOG_NAME}.uri: https://biglake.googleapis.com/iceberg/v1/restcatalog spark.sql.catalog.${OUTPUT_CATALOG_NAME}.io-impl: org.apache.iceberg.gcp.gcs.GCSFileIO spark.sql.catalog.${OUTPUT_CATALOG_NAME}.header.x-goog-user-project: ${PROJECT_ID} spark.sql.catalog.${OUTPUT_CATALOG_NAME}.warehouse: bl://projects/${PROJECT_ID}/catalogs/${OUTPUT_CATALOG_NAME} spark.sql.catalog.${OUTPUT_CATALOG_NAME}.rest.auth.type: org.apache.iceberg.gcp.auth.GoogleAuthManager spark.sql.extensions: org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions EOF属性 用途 / 说明 spark.sql.catalog.${OUTPUT_CATALOG_NAME}为自定义目录名称注册 Apache Iceberg 的目录插件 ( org.apache.iceberg.spark.SparkCatalog)。type将目录类型设置为使用 Apache Iceberg REST 目录规范 ( rest)。uriBigLake REST Metastore 的 REST API 端点网址。 io-impl将 Iceberg 配置为使用 Cloud Storage FileIO ( org.apache.iceberg.gcp.gcs.GCSFileIO) 来执行高性能数据和元数据读/写操作。header.x-goog-user-project将您的 Google Cloud 项目 ID 作为请求标头传递,以进行 BigLake API 配额和结算归因。 warehouse指向 Google Cloud 项目中的特定目录资源的 BigLake 资源 URI ( bl://projects/...)。rest.auth.type使用 Google Cloud 凭据 ( org.apache.iceberg.gcp.auth.GoogleAuthManager) 自动对 REST 目录 API 调用进行身份验证。spark.sql.extensions在 Spark SQL 中启用 Iceberg SQL 扩展和 DataFrameWriterV2(writeTo) 支持。提交 PySpark 批量作业。
设置运行时版本和目标表变量,然后使用
gcloud dataproc batches submit pyspark提交批量作业:export RUNTIME_VERSION="3.0" export SCRIPT_FILE="taxi_to_iceberg.py" gcloud dataproc batches submit pyspark "${SCRIPT_FILE}" \ --flags-file=iceberg-flags.yaml \ --project="${PROJECT_ID}" \ --region="${REGION}" \ --version="${RUNTIME_VERSION}" \ --subnet="${SUBNET_NAME}" \ --deps-bucket="gs://${BUCKET_NAME}" \ -- \ "${INPUT_PARQUET_PATH}" \ "${OUTPUT_CATALOG_NAME}" \ "${OUTPUT_DATASET_NAME}" \ "${OUTPUT_TABLE_NAME}"注意:
- 提交本地 Python 脚本时,gcloud CLI 会使用
--deps-bucket将脚本文件暂存到 Cloud Storage 中,然后再启动作业。 - 首次使用最终用户凭据 (EUC) 提交运行时
3.0时,您会收到一次性 OAuth 权限请求提示。授予访问权限,然后重新提交。 - 使用运行时
3.0时,您可能会在执行程序拆解期间在驱动程序日志中看到来自DataprocRMExecutorsAllocator的警告。假设批次达到SUCCEEDED,这些警告是暂时性的,不会导致致命错误。
- 提交本地 Python 脚本时,gcloud CLI 会使用
第 5 步:验证并查询 Iceberg 表
验证 Cloud Storage 中的输出文件,确认 BigLake 中的表注册,并使用 BigQuery 查询该表。
检查 Cloud Storage 中的 Iceberg 文件。
验证 Iceberg 数据仓库是否包含预期的
metadata/和data/目录:gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"验证 BigLake 表注册。
gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \ --catalog="${OUTPUT_CATALOG_NAME}" \ --namespace="${OUTPUT_DATASET_NAME}" \ --project="${PROJECT_ID}"从 BigQuery 查询。
BigLake 会自动在 BigQuery 中注册 Iceberg 表:
bq query \ --project_id="${PROJECT_ID}" \ --location="${REGION}" \ --use_legacy_sql=false \ "SELECT trip_date, payment_type, total_trips, avg_fare, avg_tip_percentage, total_revenue FROM \`${PROJECT_ID}.${OUTPUT_CATALOG_NAME}.${OUTPUT_DATASET_NAME}.${OUTPUT_TABLE_NAME}\` ORDER BY total_trips DESC LIMIT 10"
清理
为避免系统向您的 Google Cloud 账号收取费用,请删除您在本教程中创建的资源。
# 1. Delete BigLake Iceberg catalog, and table metadata
gcloud biglake iceberg catalogs delete "${OUTPUT_CATALOG_NAME}" \
--project="${PROJECT_ID}" --quiet 2>/dev/null || true
# 2. Delete BigQuery dataset (if created)
bq rm -r -f -d "${PROJECT_ID}:${OUTPUT_DATASET_NAME}" 2>/dev/null || true
# 3. Delete Cloud Storage bucket and warehouse files
gcloud storage rm -r "gs://${BUCKET_NAME}"