データを変換して Apache Iceberg テーブルに書き込む

このチュートリアルでは、Cloud Storage の Parquet ファイルからニューヨーク市(NYC)のタクシー乗車レコードを読み取り、Managed Service for Apache Spark サーバーレス ランタイム 3.0 で Apache Spark を使用して変換し、BigLake REST カタログを使用して集計結果を Apache Iceberg テーブルに書き込む方法について説明します。

このチュートリアルでは、次のタスクを行います。

  1. ニューヨーク市のタクシーの Parquet データをステージングします。
  2. PySpark ジョブ スクリプトを記述します。
  3. サーバーレス バッチジョブを送信します。
  4. Iceberg テーブルを確認してクエリします。

始める前に

プロジェクトを設定し、他のスタートアップ タスクを実行します。

Google Cloud プロジェクトを設定する

必要に応じてプロジェクトを設定して、API を有効にし、Identity and Access Management(IAM)ロールを付与し、アプリケーションのデフォルト認証情報を認証し、Cloud Storage バケットを作成します。

API を有効にする

Google Cloud コンソールを使用して、必要な API を有効にします。

  1. 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 the resourcemanager.projects.create permission. Learn how to grant roles.

    Go to project selector

  2. Verify that billing is enabled for your Google Cloud project.

  3. 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.enable permission. 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.

    Enable the APIs

  4. ユーザー アカウントのローカル認証情報を作成します。

    gcloud auth application-default login

    認証エラーが返され、外部 ID プロバイダ(IdP)を使用している場合は、 フェデレーション ID を使用して gcloud CLI にログインしていることを確認します。

  5. Cloud Storage バケットを作成します。
    gcloud storage buckets create gs://BUCKET_NAME
    BUCKET_NAME は、バケット名の要件を満たすバケット名に置き換えます。

必要に応じて IAM ロールを付与する

このページの例を実行するには、特定の IAM ロールが必要です。組織のポリシーによっては、これらのロールがすでに付与されている場合があります。ロールの付与を確認するには、ロールを付与する必要がありますか?をご覧ください。

ロールの付与については、プロジェクト、フォルダ、組織に対するアクセス権の管理をご覧ください。

ユーザーロール

デフォルトでは、Managed Service for Apache Spark サーバーレス ランタイム 3.0 はエンドユーザー認証情報(EUC)で実行されます。サービス アカウントのロールは必要ありません。詳細については、ペルソナとサーバーレス IAM ロールをご覧ください。

サーバーレス バッチ ワークロードを送信するために必要な権限を取得するには、次の IAM ロールを付与するよう管理者に依頼してください。

管理者は、次の 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"

次のように置き換えます。

ステップ 1. BigLake Iceberg カタログを作成する

BigLake Iceberg REST Catalog を使用すると、Spark ワークロードと BigQuery で Iceberg テーブルの検出、読み取り、書き込みを行うことができます。

  1. 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
    
  2. カタログの健全性を確認します。

    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 カタログに書き込む場合は、常に writeToDataFrameWriterV2)API を使用し、最初に現在のカタログを設定して、名前空間が存在することを確認し、SQL と writeTo 呼び出しでカタログ名をバッククォートで囲みます。

ステップ 4. サーバーレス バッチジョブを送信する

構成ファイルでカタログ プロパティを定義し、バッチジョブを Managed Service for Apache Spark サーバーレスに送信します。

  1. 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)を登録します。
    type Apache Iceberg REST カタログ仕様(rest)を使用するようにカタログ タイプを設定します。
    uri BigLake REST Metastore の REST API エンドポイント URL。
    io-impl 高パフォーマンスのデータとメタデータの読み取り/書き込みオペレーションに Cloud Storage FileIO(org.apache.iceberg.gcp.gcs.GCSFileIO)を使用するように Iceberg を構成します。
    header.x-goog-user-project BigLake API の割り当てと課金の属性設定のために、 Google Cloud プロジェクト ID をリクエスト ヘッダーとして渡します。
    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 拡張機能と DataFrameWriterV2writeTo)のサポートを有効にします。
  2. 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 に達すると致命的ではなくなります。

ステップ 5: Iceberg テーブルを検証してクエリを実行する

Cloud Storage の出力ファイルを確認し、BigLake でテーブル登録を確認して、BigQuery を使用してテーブルをクエリします。

  1. Cloud Storage 内の Iceberg ファイルを検査します。

    Iceberg ウェアハウスに想定される metadata/ ディレクトリと data/ ディレクトリが含まれていることを確認します。

    gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"
    
  2. BigLake でテーブル登録を確認します。

    gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \
      --catalog="${OUTPUT_CATALOG_NAME}" \
      --namespace="${OUTPUT_DATASET_NAME}" \
      --project="${PROJECT_ID}"
    
  3. 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}"

次のステップ