שינוי נתונים וכתיבה לטבלת Apache Iceberg

במדריך הזה נסביר איך לקרוא רשומות של נסיעות במונית בניו יורק מקובץ Parquet ב-Cloud Storage, לבצע בהן טרנספורמציה באמצעות Apache Spark בסביבת זמן ריצה ללא שרת (serverless) של Managed Service for Apache Spark 3.0, ולכתוב את התוצאות המצטברות בטבלת Apache Iceberg באמצעות קטלוג REST של BigLake.

במדריך הזה תבצעו את המשימות הבאות:

  1. הכנת נתוני Parquet של מוניות בניו יורק.
  2. כותבים את סקריפט העבודה של PySpark.
  3. שולחים את משימת האצווה בלי שרת.
  4. מאמתים את טבלת Iceberg ומריצים עליה שאילתה.

לפני שמתחילים

מגדירים את הפרויקט ומבצעים משימות אחרות שקשורות להפעלה.

הגדרת Google Cloud הפרויקט

מגדירים את הפרויקט לפי הצורך כדי להפעיל ממשקי API, להעניק תפקידים ב-Identity and Access Management (IAM), לאמת את Application Default Credentials וליצור קטגוריה של 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

    אם מוחזרת שגיאת אימות ואתם משתמשים בספק זהויות חיצוני (IdP), ודאו ש נכנסתם ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.

  5. יוצרים קטגוריה של Cloud Storage:
    gcloud storage buckets create gs://BUCKET_NAME
    מחליפים את BUCKET_NAME בשם קטגוריה שעומד בקריטריונים לשמות של קטגוריות.

הקצאת תפקידים ב-IAM אם צריך

כדי להריץ את הדוגמאות בדף הזה, צריך תפקידים מסוימים ב-IAM. יכול להיות שהתפקידים האלה כבר הוקצו, בהתאם למדיניות הארגון. כדי לבדוק את התפקידים שהוקצו, ראו האם צריך להקצות תפקידים?.

מידע נוסף על הקצאת תפקידים מופיע במאמר ניהול הגישה לפרויקטים, לתיקיות ולארגונים.

תפקידי משתמשים

כברירת מחדל, זמן הריצה של Managed Service for Apache Spark serverless‏ 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

קטלוג ה-REST של BigLake Iceberg מאפשר לעומסי עבודה של Spark ול-BigQuery לגלות, לקרוא ולכתוב טבלאות Iceberg.

  1. יוצרים את קטלוג 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
    
  2. אימות תקינות הקטלוג.

    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.

  1. מגדירים את מאפייני קטלוג 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.
  2. שולחים את עבודת ה-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.

שלב 5: מאמתים את טבלת Iceberg ומריצים עליה שאילתה

בודקים את קובצי הפלט ב-Cloud Storage, מוודאים שהטבלה רשומה ב-BigLake ומריצים שאילתה על הטבלה באמצעות BigQuery.

  1. בודקים קובצי Iceberg ב-Cloud Storage.

    מוודאים שמאגר 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}"

המאמרים הבאים