Mengubah data dan menulis ke tabel Apache Iceberg

Tutorial ini menunjukkan cara membaca data perjalanan taksi New York City (NYC) dari file Parquet di Cloud Storage, mentransformasinya dengan Apache Spark di runtime serverless Managed Service untuk Apache Spark 3.0, dan menulis hasil gabungan ke dalam tabel Apache Iceberg menggunakan Katalog REST BigLake.

Dalam tutorial ini, Anda akan menyelesaikan tugas berikut:

  1. Siapkan data Parquet taksi NYC.
  2. Tulis skrip tugas PySpark.
  3. Kirimkan tugas batch serverless.
  4. Verifikasi dan kueri tabel Iceberg.

Sebelum memulai

Siapkan project Anda dan lakukan tugas startup lainnya.

Menyiapkan project Google Cloud

Siapkan project Anda sesuai kebutuhan untuk mengaktifkan API, memberikan peran Identity and Access Management (IAM), mengautentikasi Kredensial Default Aplikasi, dan membuat bucket Cloud Storage.

Mengaktifkan API

Gunakan konsol Google Cloud untuk mengaktifkan API yang diperlukan.

  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. Buat kredensial autentikasi lokal untuk akun pengguna Anda:

    gcloud auth application-default login

    Jika error autentikasi ditampilkan, dan Anda menggunakan penyedia identitas (IdP) eksternal, konfirmasi bahwa Anda telah login ke gcloud CLI dengan identitas gabungan Anda.

  5. Membuat bucket Cloud Storage:
    gcloud storage buckets create gs://BUCKET_NAME
    Ganti BUCKET_NAME dengan nama bucket yang memenuhi persyaratan penamaan bucket.

Memberikan peran IAM jika diperlukan

Peran IAM tertentu diperlukan untuk menjalankan contoh di halaman ini. Bergantung pada kebijakan organisasi, peran ini mungkin sudah diberikan. Untuk memeriksa pemberian peran, lihat Apakah Anda perlu memberikan peran?.

Untuk mengetahui informasi selengkapnya tentang cara memberikan peran, lihat Mengelola akses ke project,folder, dan organisasi.

Peran pengguna

Secara default, runtime serverless Managed Service untuk Apache Spark 3.0 berjalan dengan kredensial pengguna akhir (EUC) Anda. Peran akun layanan tidak diperlukan. Untuk mengetahui informasi selengkapnya, lihat Persona dan peran IAM serverless.

Untuk mendapatkan izin yang Anda perlukan untuk mengirimkan workload batch serverless, minta administrator untuk memberi Anda peran IAM berikut:

Administrator Anda dapat menjalankan skrip bash berikut untuk memberikan peran ke akun pengguna Anda.

```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
```

Mengonfigurasi variabel lingkungan

Jalankan skrip bash berikut untuk menetapkan variabel lingkungan shell yang digunakan di seluruh tutorial ini.

# 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"

Ganti kode berikut:

Langkah 1. Membuat katalog Iceberg BigLake

Katalog REST Iceberg BigLake memungkinkan workload Spark dan BigQuery menemukan, membaca, dan menulis tabel Iceberg.

  1. Buat katalog BigLake Iceberg.

    Buat katalog dengan lokasi defaultnya yang mengarah ke jalur gudang Anda di 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. Verifikasi kesehatan katalog.

    gcloud biglake iceberg catalogs describe "${OUTPUT_CATALOG_NAME}" --project="${PROJECT_ID}"
    

Langkah 2. Menyiapkan data Parquet taksi NYC

NYC Taxi & Limousine Commission (TLC) memublikasikan data perjalanan bulanan dalam Parquet. Download data selama satu bulan dan salin ke bucket Cloud Storage Anda.

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}"

Langkah 3. Tulis skrip tugas PySpark

Buat taxi_to_iceberg.py. Skrip ini membaca rekaman Parquet mentah, memfilter perjalanan yang tidak valid, menghitung agregat penumpang, tarif, tip, dan pendapatan harian, serta menulis tabel Iceberg menggunakan katalog REST BigLake.

"""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()

Penting: Saat menulis ke katalog BigLake Iceberg, selalu gunakan API writeTo (DataFrameWriterV2), tetapkan katalog saat ini terlebih dahulu, pastikan namespace ada, dan sertakan nama katalog dengan tanda petik terbalik dalam panggilan SQL dan writeTo.

Langkah 4. Kirimkan tugas batch serverless

Tentukan properti katalog dalam file konfigurasi dan kirimkan tugas batch ke Managed Service untuk Apache Spark serverless.

  1. Tentukan properti katalog Iceberg dalam file tanda YAML.

    Buat iceberg-flags.yaml yang berisi properti Katalog REST BigLake. Google Cloud CLI membaca file ini menggunakan argumen --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
    
    Properti Tujuan / Deskripsi
    spark.sql.catalog.${OUTPUT_CATALOG_NAME} Mendaftarkan plugin katalog Apache Iceberg (org.apache.iceberg.spark.SparkCatalog) untuk nama katalog kustom.
    type Menetapkan jenis katalog untuk menggunakan spesifikasi katalog REST Apache Iceberg (rest).
    uri URL endpoint API REST untuk BigLake REST Metastore.
    io-impl Mengonfigurasi Iceberg untuk menggunakan Cloud Storage FileIO (org.apache.iceberg.gcp.gcs.GCSFileIO) untuk operasi baca/tulis data dan metadata berperforma tinggi.
    header.x-goog-user-project Meneruskan ID project Google Cloud Anda sebagai header permintaan untuk atribusi penagihan dan kuota BigLake API.
    warehouse URI Resource BigLake (bl://projects/...) yang mengarah ke resource katalog tertentu di project Google Cloud Anda.
    rest.auth.type Mengautentikasi panggilan REST catalog API secara otomatis menggunakan Google Cloud kredensial (org.apache.iceberg.gcp.auth.GoogleAuthManager).
    spark.sql.extensions Mengaktifkan ekstensi SQL Iceberg dan dukungan DataFrameWriterV2 (writeTo) di Spark SQL.
  2. Kirimkan tugas batch PySpark.

    Tetapkan variabel versi runtime dan tabel target, lalu kirimkan tugas batch menggunakan 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}"
    

    Catatan:

    • Saat mengirimkan skrip Python lokal, gcloud CLI menggunakan --deps-bucket untuk menyiapkan file skrip ke Cloud Storage sebelum memulai tugas.
    • Pada pengiriman 3.0 runtime pertama menggunakan kredensial pengguna akhir (EUC), Anda akan menerima perintah izin OAuth satu kali. Berikan akses dan kirim ulang.
    • Saat menggunakan 3.0 runtime, Anda mungkin melihat peringatan dari DataprocRMExecutorsAllocator dalam log driver selama penonaktifan eksekutor. Peringatan ini bersifat sementara dan tidak fatal jika batch mencapai SUCCEEDED.

Langkah 5: Verifikasi dan kueri tabel Iceberg

Verifikasi file output di Cloud Storage, konfirmasi pendaftaran tabel di BigLake, dan buat kueri tabel menggunakan BigQuery.

  1. Periksa file Iceberg di Cloud Storage.

    Verifikasi bahwa gudang Iceberg berisi direktori metadata/ dan data/ yang diharapkan:

    gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"
    
  2. Verifikasi pendaftaran tabel dengan BigLake.

    gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \
      --catalog="${OUTPUT_CATALOG_NAME}" \
      --namespace="${OUTPUT_DATASET_NAME}" \
      --project="${PROJECT_ID}"
    
  3. Buat kueri dari BigQuery.

    BigLake otomatis mendaftarkan tabel Iceberg di 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"
    

Pembersihan

Agar akun Google Cloud Anda tidak dikenai biaya, hapus resource yang Anda buat dalam tutorial ini.

# 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}"

Langkah berikutnya