このチュートリアルでは、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
認証エラーが返され、外部 ID プロバイダ(IdP)を使用している場合は、 フェデレーション ID を使用して 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 ロールを付与するよう管理者に依頼してください。
-
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 バケットの名前。
ステップ 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 データをステージングする
NYC Taxi & Limousine Commission(TLC)は、月ごとの乗車記録を Parquet で公開しています。1 か月分のデータをダウンロードして、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 レコードを読み取り、無効な乗車をフィルタし、乗客、運賃、チップ、収益の 1 日あたりの集計を計算し、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 カタログのプロパティを定義します。
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)を登録します。typeApache Iceberg REST カタログ仕様( rest)を使用するようにカタログ タイプを設定します。uriBigLake REST Metastore の REST API エンドポイント URL。 io-impl高パフォーマンスのデータとメタデータの読み取り/書き込みオペレーションに Cloud Storage FileIO( org.apache.iceberg.gcp.gcs.GCSFileIO)を使用するように Iceberg を構成します。header.x-goog-user-projectBigLake API の割り当てと課金の属性設定のために、 Google Cloud プロジェクト ID をリクエスト ヘッダーとして渡します。 warehouseGoogle Cloud プロジェクト内の特定のカタログ リソースを指す BigLake リソース URI( bl://projects/...)。rest.auth.typeGoogle Cloud 認証情報( org.apache.iceberg.gcp.auth.GoogleAuthManager)を使用して REST カタログ API 呼び出しを自動的に認証します。spark.sql.extensionsSpark 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 同意プロンプトが 1 回表示されます。アクセス権を付与して、もう一度送信してください。 - ランタイム
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 は、Iceberg テーブルを BigQuery に自動的に登録します。
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 の問題のトラブルシューティングをご覧ください。