במדריך הזה נסביר איך לקרוא רשומות של נסיעות במונית בניו יורק מקובץ Parquet ב-Cloud Storage, לבצע בהן טרנספורמציה באמצעות Apache Spark בסביבת זמן ריצה ללא שרת (serverless) של Managed Service for Apache Spark 3.0, ולכתוב את התוצאות המצטברות בטבלת Apache Iceberg באמצעות קטלוג REST של BigLake.
במדריך הזה תבצעו את המשימות הבאות:
- הכנת נתוני Parquet של מוניות בניו יורק.
- כותבים את סקריפט העבודה של PySpark.
- שולחים את משימת האצווה בלי שרת.
- מאמתים את טבלת Iceberg ומריצים עליה שאילתה.
לפני שמתחילים
מגדירים את הפרויקט ומבצעים משימות אחרות שקשורות להפעלה.
הגדרת Google Cloud הפרויקט
מגדירים את הפרויקט לפי הצורך כדי להפעיל ממשקי API, להעניק תפקידים ב-Identity and Access Management (IAM), לאמת את Application Default Credentials וליצור קטגוריה של 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
אם מוחזרת שגיאת אימות ואתם משתמשים בספק זהויות חיצוני (IdP), ודאו ש נכנסתם ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.
-
יוצרים קטגוריה של Cloud Storage:
מחליפים אתgcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAMEבשם קטגוריה שעומד בקריטריונים לשמות של קטגוריות.
הקצאת תפקידים ב-IAM אם צריך
כדי להריץ את הדוגמאות בדף הזה, צריך תפקידים מסוימים ב-IAM. יכול להיות שהתפקידים האלה כבר הוקצו, בהתאם למדיניות הארגון. כדי לבדוק את התפקידים שהוקצו, ראו האם צריך להקצות תפקידים?.
מידע נוסף על הקצאת תפקידים מופיע במאמר ניהול הגישה לפרויקטים, לתיקיות ולארגונים.
תפקידי משתמשים
כברירת מחדל, זמן הריצה של Managed Service for Apache Spark serverless 3.0 פועל עם פרטי הכניסה של משתמש הקצה (EUC). לא נדרשים תפקידים של חשבון שירות. מידע נוסף זמין במאמר בנושא פרסונות ותפקידי IAM ללא שרת.
כדי לקבל את ההרשאות שדרושות לשליחת עומס עבודה של אצווה ללא שרתים, צריך לבקש מהאדמין להקצות לכם את תפקידי ה-IAM הבאים:
-
הפעלת עומסי עבודה ב-Runtime 3.x (ברירת מחדל של EUC):
- עורך של Dataproc Serverless (
roles/dataproc.serverlessEditor) בפרויקט - אדמין BigLake (
roles/biglake.admin) בפרויקט - BigQuery Admin (
roles/bigquery.admin) on the project - אדמין של אובייקט אחסון (
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: השם של קטגוריית Cloud Storage שיצרתם בקטע הגדרת הפרויקט Google Cloud .
שלב 1. יצירת קטלוג BigLake Iceberg
קטלוג ה-REST של BigLake Iceberg מאפשר לעומסי עבודה של Spark ול-BigQuery לגלות, לקרוא ולכתוב טבלאות Iceberg.
יוצרים את קטלוג Iceberg של BigLake.
יוצרים את הקטלוג עם מיקום ברירת המחדל שלו שמפנה לנתיב של מחסן הנתונים ב-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.
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 הגולמיות, מסנן נסיעות לא תקינות, מחשב את נתוני הנוסעים, התעריפים, הטיפים וההכנסות המצטברים היומיים, וכותב טבלת Iceberg באמצעות קטלוג BigLake REST.
"""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, צריך תמיד להשתמש ב-API writeTo (DataFrameWriterV2), להגדיר קודם את הקטלוג הנוכחי, לוודא שמרחב השמות קיים ולהוסיף גרשיים הפוכים לשם הקטלוג ב-SQL ובקריאות ל-API writeTo.
שלב 4. שליחת משימת אצווה בלי שרת
מגדירים את מאפייני הקטלוג בקובץ הגדרות ושולחים את משימת האצווה ל-Managed Service for Apache Spark serverless.
מגדירים את מאפייני קטלוג Iceberg בקובץ flags בפורמט YAML.
יוצרים
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מאפיין (property) מטרה / תיאור spark.sql.catalog.${OUTPUT_CATALOG_NAME}רושם את תוסף הקטלוג של Apache Iceberg ( org.apache.iceberg.spark.SparkCatalog) עבור שם הקטלוג המותאם אישית.typeמגדירים את סוג הקטלוג לשימוש במפרט של קטלוג REST של Apache Iceberg ( rest).uriכתובת ה-URL של נקודת הקצה של BigLake REST Metastore. io-implההגדרה קובעת ש-Iceberg ישתמש ב-Cloud Storage FileIO ( org.apache.iceberg.gcp.gcs.GCSFileIO) לפעולות קריאה וכתיבה של נתונים ומטא-נתונים עם ביצועים גבוהים.header.x-goog-user-projectמעביר את מזהה הפרויקט Google Cloud ככותרת בקשה למכסת BigLake API ולשיוך חיובים. warehouseה-URI של משאב BigLake ( bl://projects/...) שמצביע על משאב הקטלוג הספציפי בפרויקט Google Cloud .rest.auth.typeאימות אוטומטי של קריאות API לקטלוג REST באמצעות Google Cloud פרטי כניסה ( org.apache.iceberg.gcp.auth.GoogleAuthManager).spark.sql.extensionsהאפשרות הזו מפעילה תוספים של Iceberg SQL ותמיכה ב- DataFrameWriterV2(writeTo) ב-Spark SQL.שולחים את עבודת ה-batch של 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 מקומי, ה-CLI של gcloud משתמש ב-
--deps-bucketכדי להכין את קובץ הסקריפט ב-Cloud Storage לפני הפעלת העבודה. - בפעם הראשונה שתשלחו את
3.0בזמן ריצה באמצעות פרטי הכניסה של משתמש הקצה (EUC), תופיע בקשה חד-פעמית להסכמה ל-OAuth. צריך לתת גישה ולשלוח מחדש. - כשמשתמשים ב-
3.0בזמן ריצה, יכול להיות שיופיעו אזהרות מ-DataprocRMExecutorsAllocatorביומני מנהלי ההתקנים במהלך פירוק של מנהל הביצוע. האזהרות האלה הן זמניות ולא קריטיות, בהנחה שהקבוצה מגיעה אלSUCCEEDED.
- כששולחים סקריפט Python מקומי, ה-CLI של gcloud משתמש ב-
שלב 5: מאמתים את טבלת Iceberg ומריצים עליה שאילתה
בודקים את קובצי הפלט ב-Cloud Storage, מוודאים שהטבלה רשומה ב-BigLake ומריצים שאילתה על הטבלה באמצעות BigQuery.
בודקים קובצי Iceberg ב-Cloud Storage.
מוודאים שמאגר 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}"