本教學課程說明如何從 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 bucket。
啟用 API
使用 Google Cloud 控制台啟用必要的 API。
- 登入 Google Cloud 帳戶。如果您是 Google Cloud新手,歡迎 建立帳戶,親自評估產品在實際工作環境中的成效。新客戶還能獲得價值 $300 美元的免費抵免額,可用於執行、測試及部署工作負載。
-
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, if any are not already enabled.
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.-
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, if any are not already enabled.
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
如果您使用 Cloud Shell,則不需要執行這項操作。
如果系統傳回驗證錯誤,且您使用外部識別資訊提供者 (IdP),請確認您已 使用聯合身分登入 gcloud CLI。
-
建立 Cloud Storage bucket:
將gcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAME替換為符合值區命名規定的值區名稱。
視需要授予 IAM 角色
如要執行本頁的範例,您必須具備特定 IAM 角色。 根據機構政策,這些角色可能已獲授權。 如要檢查角色授予情形,請參閱「是否需要授予角色?」一文。
如要進一步瞭解如何授予角色,請參閱「管理專案、資料夾和機構的存取權」。
使用者角色
根據預設,Managed Service for Apache Spark 無伺服器執行階段 3.0 會在您的使用者憑證 (EUC) 下執行。不需要服務帳戶角色。詳情請參閱「角色和無伺服器 IAM 角色」一文。
如要取得提交無伺服器批次工作負載所需的權限,請要求管理員授予您下列 IAM 角色:
-
在 Runtime 3.x (預設 EUC) 上執行工作負載:
- 專案的 Dataproc Serverless 編輯者 (
roles/dataproc.serverlessEditor) - 專案的 BigLake 管理員 (
roles/biglake.admin) - 專案的 BigQuery 管理員 (
roles/bigquery.admin) - 專案的Storage 物件管理員 (
roles/storage.objectAdmin)
- 專案的 Dataproc Serverless 編輯者 (
管理員可以執行下列 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 指令碼,設定本教學課程中使用的殼層環境變數。
# 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 bucket 名稱。
步驟 1:建立 BigLake Iceberg 目錄
BigLake Iceberg REST 目錄可讓 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 bucket。
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 目錄屬性。
建立
iceberg-flags.yaml,其中包含 BigLake REST 目錄屬性。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}"
後續步驟
- 進一步瞭解 Managed Service for Apache Spark 無伺服器批次工作負載。
- 探索角色和無伺服器 IAM 角色。
- 請參閱 BigQuery 中的 BigLake Iceberg 資料表。
- 請參閱排解常見的 Lakehouse 問題。