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:
- Siapkan data Parquet taksi NYC.
- Tulis skrip tugas PySpark.
- Kirimkan tugas batch serverless.
- 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.
-
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.-
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.
-
Membuat bucket Cloud Storage:
Gantigcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAMEdengan 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:
-
Menjalankan workload di Runtime 3.x (EUC default):
- Dataproc Serverless Editor (
roles/dataproc.serverlessEditor) di project - Admin BigLake (
roles/biglake.admin) di project - Admin BigQuery (
roles/bigquery.admin) di project - Storage Object Admin (
roles/storage.objectAdmin) di project
- Dataproc Serverless Editor (
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:
BUCKET_NAME: nama bucket Cloud Storage yang Anda buat di Siapkan project Google Cloud .
Langkah 1. Membuat katalog Iceberg BigLake
Katalog REST Iceberg BigLake memungkinkan workload Spark dan BigQuery menemukan, membaca, dan menulis tabel Iceberg.
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-userVerifikasi 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.
Tentukan properti katalog Iceberg dalam file tanda YAML.
Buat
iceberg-flags.yamlyang 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 EOFProperti Tujuan / Deskripsi spark.sql.catalog.${OUTPUT_CATALOG_NAME}Mendaftarkan plugin katalog Apache Iceberg ( org.apache.iceberg.spark.SparkCatalog) untuk nama katalog kustom.typeMenetapkan jenis katalog untuk menggunakan spesifikasi katalog REST Apache Iceberg ( rest).uriURL endpoint API REST untuk BigLake REST Metastore. io-implMengonfigurasi 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-projectMeneruskan ID project Google Cloud Anda sebagai header permintaan untuk atribusi penagihan dan kuota BigLake API. warehouseURI Resource BigLake ( bl://projects/...) yang mengarah ke resource katalog tertentu di project Google Cloud Anda.rest.auth.typeMengautentikasi panggilan REST catalog API secara otomatis menggunakan Google Cloud kredensial ( org.apache.iceberg.gcp.auth.GoogleAuthManager).spark.sql.extensionsMengaktifkan ekstensi SQL Iceberg dan dukungan DataFrameWriterV2(writeTo) di Spark SQL.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-bucketuntuk menyiapkan file skrip ke Cloud Storage sebelum memulai tugas. - Pada pengiriman
3.0runtime pertama menggunakan kredensial pengguna akhir (EUC), Anda akan menerima perintah izin OAuth satu kali. Berikan akses dan kirim ulang. - Saat menggunakan
3.0runtime, Anda mungkin melihat peringatan dariDataprocRMExecutorsAllocatordalam log driver selama penonaktifan eksekutor. Peringatan ini bersifat sementara dan tidak fatal jika batch mencapaiSUCCEEDED.
- Saat mengirimkan skrip Python lokal, gcloud CLI menggunakan
Langkah 5: Verifikasi dan kueri tabel Iceberg
Verifikasi file output di Cloud Storage, konfirmasi pendaftaran tabel di BigLake, dan buat kueri tabel menggunakan BigQuery.
Periksa file Iceberg di Cloud Storage.
Verifikasi bahwa gudang Iceberg berisi direktori
metadata/dandata/yang diharapkan:gcloud storage ls -r "gs://${BUCKET_NAME}/warehouse/${OUTPUT_DATASET_NAME}/${OUTPUT_TABLE_NAME}/"Verifikasi pendaftaran tabel dengan BigLake.
gcloud biglake iceberg tables describe "${OUTPUT_TABLE_NAME}" \ --catalog="${OUTPUT_CATALOG_NAME}" \ --namespace="${OUTPUT_DATASET_NAME}" \ --project="${PROJECT_ID}"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
- Pelajari lebih lanjut workload batch serverless Managed Service untuk Apache Spark.
- Pelajari Persona dan peran IAM serverless.
- Baca tabel Iceberg BigLake di BigQuery.
- Lihat Memecahkan masalah umum Lakehouse.