שליחת שאילתות לטבלאות Iceberg באמצעות קטלוג זמן הריצה של Lakehouse, ‏ Spark ו-BigQuery

במאמר הזה נסביר איך ליצור קטלוג זמן ריצה של Lakehouse עם קטלוג של כמה דליים כדי להשתמש ב-Lakehouse for Apache Iceberg. ההגדרה הזו יוצרת שכבת מטא-נתונים מנוהלת שמקשרת בין מנועי עיבוד קוד פתוח לבין Google Cloud.

לאחר מכן מריצים משימת PySpark ב-Managed Service for Apache Spark כדי ליצור טבלה של Lakehouse Iceberg REST catalog באמצעות נקודת הקצה של Apache Iceberg REST catalog.

לאחר מכן, תוכלו להריץ שאילתות על הטבלה שמתקבלת ישירות ממסוף Google Cloud ב-BigQuery באמצעות התחביר של project.catalog.namespace.table.

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

  1. נכנסים לחשבון Google Cloud . אם אתם משתמשים חדשים ב- Google Cloud, צרו חשבון כדי שתוכלו להעריך את הביצועים של המוצרים שלנו בתרחישים מהעולם האמיתי. לקוחות חדשים מקבלים בחינם גם קרדיט בשווי 300$ להרצה, לבדיקה ולפריסה של עומסי העבודה.
  2. 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

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

  4. Enable the BigLake, Dataproc 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

  5. 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

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

  7. Enable the BigLake, Dataproc 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

מתן תפקידים ב-IAM

כדי לאפשר למשימת PySpark של Managed Service for Apache Spark ולספריית זמן הריצה של Lakehouse לעבוד עם Cloud Storage ו-BigQuery, צריך להעניק את התפקידים הנדרשים לחשבונות המשתמשים המתאימים:

  1. במסוף Google Cloud , לוחצים על הפעלת Cloud Shell.

    הפעלת Cloud Shell

  2. לוחצים על Authorize.

  3. מעניקים את התפקיד Dataproc Worker לחשבון השירות שמשמש כברירת המחדל של Compute Engine בפרויקט, שמשמש כברירת מחדל את Managed Service for Apache Spark, כפי שמפורט במאמר חשבונות שירות של Managed Service for Apache Spark.

    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:$(gcloud projects describe PROJECT_ID --format='value(projectNumber)')-compute@developer.gserviceaccount.com" \
        --role="roles/dataproc.worker"
  4. מקצים את התפקיד Service Usage Consumer לחשבון השירות שמשמש כברירת המחדל של Compute Engine בפרויקט.

    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:$(gcloud projects describe PROJECT_ID --format='value(projectNumber)')-compute@developer.gserviceaccount.com" \
        --role="roles/serviceusage.serviceUsageConsumer"
  5. מקצים את התפקיד BigLake Editor לחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine של הפרויקט.

    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:$(gcloud projects describe PROJECT_ID --format='value(projectNumber)')-compute@developer.gserviceaccount.com" \
        --role="roles/biglake.editor"
  6. מעניקים את התפקיד BigQuery Data Editor לחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine של הפרויקט.

    gcloud projects add-iam-policy-binding PROJECT_ID \
        --member="serviceAccount:$(gcloud projects describe PROJECT_ID --format='value(projectNumber)')-compute@developer.gserviceaccount.com" \
        --role="roles/bigquery.dataEditor"

    מחליפים את מה שכתוב בשדות הבאים:

    • PROJECT_ID: מזהה הפרויקט ב- Google Cloud

יצירת קטלוג של Lakehouse בזמן ריצה

יוצרים קטלוג של זמן ריצה של Lakehouse כדי לנהל את המטא-נתונים של טבלאות Iceberg.

  1. ב-Cloud Shell, מריצים את הפקודה הבאה כדי ליצור את הקטגוריה עם כמה קטגוריות (bl://) עם מכירת אישורים:

    gcloud biglake iceberg catalogs create LAKEHOUSE_CATALOG_ID \
        --project=PROJECT_ID \
        --catalog-type=biglake \
        --default-location=gs://BUCKET_NAME \
        --credential-mode=vended-credentials

    מחליפים את מה שכתוב בשדות הבאים:

    • LAKEHOUSE_CATALOG_ID: שם ייחודי לקטלוג.
    • PROJECT_ID: מזהה הפרויקט ב- Google Cloud .
    • BUCKET_NAME: השם של קטגוריית Cloud Storage שמכילה את קובץ האפליקציה של PySpark.
  2. נותנים הרשאות לחשבון השירות של הקטלוג בדלי:

    1. במסוף Google Cloud , עוברים אל Lakehouse.

      מעבר אל Lakehouse

    2. לוחצים על השם של הקטלוג שיצרתם (LAKEHOUSE_CATALOG_ID).

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

    4. בתיבת הדו-שיח, לוחצים על אישור. כך מוודאים שלחשבון השירות של הקטלוג יש את התפקיד Storage Object User בדלי.

יצירה והרצה של משימת PySpark

כדי ליצור טבלת Iceberg ולשאול עליה שאילתות, קודם צריך ליצור משימת PySpark עם הצהרות Spark SQL הנדרשות. לאחר מכן מריצים את המשימה באמצעות Managed Service for Apache Spark.

יצירת סקריפט PySpark עם מרחב שמות וטבלה

בעורך טקסט, יוצרים קובץ בשם quickstart.py עם התוכן הבא.

סקריפט PySpark הזה מאתחל סשן Spark כדי לבצע כמה פעולות בקטלוג Iceberg. הסקריפט יוצר תחילה מרחב שמות, אם הוא עדיין לא קיים. לאחר מכן נוצרת טבלת Iceberg בשם quickstart_table עם סכימה בסיסית. אחרי שהטבלה נוצרת, הסקריפט מוסיף שלוש שורות של נתונים. לבסוף, היא שולחת שאילתה לטבלה כדי לאחזר את כל הרשומות שהוכנסו.

אחר כך משתמשים בערכים האלה בשלב הבא כשמריצים את העבודה gcloud dataproc batches submit pyspark.

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("quickstart").getOrCreate()

# Create a namespace (dataset) if it doesn't exist
spark.sql("CREATE NAMESPACE IF NOT EXISTS `quickstart_catalog`.quickstart_namespace")

# Create the table
spark.sql("""
    CREATE OR REPLACE TABLE `quickstart_catalog`.quickstart_namespace.quickstart_table (
        id INT,
        name STRING
    )
    USING iceberg
""")

# Insert data into the table
spark.sql("""
    INSERT INTO `quickstart_catalog`.quickstart_namespace.quickstart_table
    VALUES (1, 'one'), (2, 'two'), (3, 'three')
""")

העלאת הסקריפט לקטגוריה של Cloud Storage

אחרי שיוצרים את סקריפט quickstart.py, מעלים אותו לקטגוריה ב-Cloud Storage.

  1. במסוף Google Cloud , עוברים אל Cloud Storage buckets.

    כניסה לדף Buckets

  2. לוחצים על שם הקטגוריה.

  3. בכרטיסייה Objects, לוחצים על Upload > Upload files.

  4. בדפדפן הקבצים, בוחרים את קובץ quickstart.py ולוחצים על פתיחה.

הפעלת עבודת PySpark

אחרי שמעלים את סקריפט quickstart.py, מריצים אותו כעבודת עיבוד ברצף (batch job) של Managed Service for Apache Spark.

  1. ב-Cloud Shell, מריצים את משימת האצווה של Managed Service for Apache Spark באמצעות הסקריפט quickstart.py.

    gcloud dataproc batches submit pyspark gs://BUCKET_NAME/quickstart.py \
        --project=PROJECT_ID \
        --region=REGION \
        --version=2.2 \
        --properties="\
    spark.sql.defaultCatalog=quickstart_catalog,\
    spark.sql.catalog.quickstart_catalog=org.apache.iceberg.spark.SparkCatalog,\
    spark.sql.catalog.quickstart_catalog.type=rest,\
    spark.sql.catalog.quickstart_catalog.uri=https://biglake.googleapis.com/iceberg/v1/restcatalog,\
    spark.sql.catalog.quickstart_catalog.warehouse=bl://projects/PROJECT_ID/catalogs/LAKEHOUSE_CATALOG_ID,\
    spark.sql.catalog.quickstart_catalog.io-impl=org.apache.iceberg.gcp.gcs.GCSFileIO,\
    spark.sql.catalog.quickstart_catalog.header.x-goog-user-project=PROJECT_ID,\
    spark.sql.catalog.quickstart_catalog.rest.auth.type=org.apache.iceberg.gcp.auth.GoogleAuthManager,\
    spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions,\
    spark.sql.catalog.quickstart_catalog.header.X-Iceberg-Access-Delegation=vended-credentials,\
    spark.sql.catalog.quickstart_catalog.gcs.oauth2.refresh-credentials-endpoint=https://oauth2.googleapis.com/token"

    מחליפים את מה שכתוב בשדות הבאים:

    • BUCKET_NAME: השם של קטגוריית Cloud Storage שמכילה את קובץ האפליקציה של PySpark.

    • LAKEHOUSE_CATALOG_ID: השם של קטלוג BigLake. השם הזה ישמש בהמשך כשמבצעים שאילתה בקטלוג ב-BigQuery, באמצעות התחביר של P.C.N.T. לדוגמה, my-project.biglake-catalog.quickstart_namespace.quickstart_table.

    • PROJECT_ID: מזהה הפרויקט ב- Google Cloud .

    • REGION: האזור שבו יפעל עומס העבודה של עיבוד ברצף (batch processing) ב-Managed Service for Apache Spark.

    כשהפעולה מסתיימת, מוצג פלט שדומה לזה:

    Batch [cb9d84e9489d408baca4f9e7ab4c64ff] finished.
    metadata:
    '@type': type.googleapis.com/google.cloud.dataproc.v1.BatchOperationMetadata
    batch: projects/your-project/locations/us-central1/batches/cb9d84e9489d408baca4f9e7ab4c64ff
    batchUuid: 54b0b9d2-f0a1-4fdf-ae44-eead3f8e60e9
    createTime: '2026-01-24T00:10:50.224097Z'
    description: Batch
    labels:
        goog-dataproc-batch-id: cb9d84e9489d408baca4f9e7ab4c64ff
        goog-dataproc-batch-uuid: 54b0b9d2-f0a1-4fdf-ae44-eead3f8e60e9
        goog-dataproc-drz-resource-uuid: batch-54b0b9d2-f0a1-4fdf-ae44-eead3f8e60e9
        goog-dataproc-location: us-central1
    operationType: BATCH
    name: projects/your-project/regions/us-central1/operations/32287926-5f61-3572-b54a-fbad8940d6ef
    

שליחת שאילתה לטבלה מ-BigQuery

  1. במסוף Google Cloud , עוברים אל BigQuery.

    כניסה ל-BigQuery

  2. מזינים את ההצהרה הבאה בעורך השאילתות. השאילתה משתמשת בתחביר project.catalog.namespace.table.

    SELECT * FROM `PROJECT_ID.LAKEHOUSE_CATALOG_ID.quickstart_namespace.quickstart_table`;
    

    מחליפים את:

    • PROJECT_ID: מזהה הפרויקט ב- Google Cloud .

    • LAKEHOUSE_CATALOG_ID: המזהה של הקטלוג לשימוש בשאילתות BigQuery.

  3. לוחצים על Run.

    בתוצאות השאילתה מוצגים הנתונים שהוספתם באמצעות עבודת ה-PySpark.

הסרת המשאבים

כדי לא לצבור חיובים לחשבון Google Cloud על המשאבים שבהם השתמשתם בדף הזה, פועלים לפי השלבים הבאים:

  1. כדי למחוק את מרחב השמות (קבוצת הנתונים) והטבלה, מעדכנים את quickstart.py:

    from pyspark.sql import SparkSession
    
    spark = SparkSession.builder.appName("quickstart").getOrCreate()
    
    # Delete the table first, then the namespace (dataset)
    spark.sql("DROP TABLE `quickstart_catalog`.quickstart_namespace.quickstart_table")
    spark.sql("DROP NAMESPACE `quickstart_catalog`.quickstart_namespace")
    

    מעלים אותו לקטגוריה של Cloud Storage:

    1. במסוף Google Cloud , עוברים אל Cloud Storage buckets.

      כניסה לדף Buckets

    2. לוחצים על שם הקטגוריה.

    3. בכרטיסייה Objects, לוחצים על Upload > Upload files.

    4. בדפדפן הקבצים, בוחרים את קובץ quickstart.py ולוחצים על פתיחה.

    ב-Cloud Shell, מריצים עוד משימה באצווה של Managed Service for Apache Spark באמצעות הסקריפט המעודכן quickstart.py.

    gcloud dataproc batches submit pyspark gs://BUCKET_NAME/quickstart.py \
        --project=PROJECT_ID \
        --region=REGION \
        --version=2.2 \
        --properties="\
    spark.sql.defaultCatalog=quickstart_catalog,\
    spark.sql.catalog.quickstart_catalog=org.apache.iceberg.spark.SparkCatalog,\
    spark.sql.catalog.quickstart_catalog.type=rest,\
    spark.sql.catalog.quickstart_catalog.uri=https://biglake.googleapis.com/iceberg/v1/restcatalog,\
    spark.sql.catalog.quickstart_catalog.warehouse=bl://projects/PROJECT_ID/catalogs/LAKEHOUSE_CATALOG_ID,\
    spark.sql.catalog.quickstart_catalog.io-impl=org.apache.iceberg.gcp.gcs.GCSFileIO,\
    spark.sql.catalog.quickstart_catalog.header.x-goog-user-project=PROJECT_ID,\
    spark.sql.catalog.quickstart_catalog.rest.auth.type=org.apache.iceberg.gcp.auth.GoogleAuthManager,\
    spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions,\
    spark.sql.catalog.quickstart_catalog.header.X-Iceberg-Access-Delegation=vended-credentials,\
    spark.sql.catalog.quickstart_catalog.gcs.oauth2.refresh-credentials-endpoint=https://oauth2.googleapis.com/token"

    מחליפים את מה שכתוב בשדות הבאים:

    • BUCKET_NAME: השם של קטגוריית Cloud Storage שמכילה את קובץ האפליקציה של PySpark.

    • LAKEHOUSE_CATALOG_ID: השם של קטלוג BigLake.

    • PROJECT_ID: מזהה הפרויקט ב- Google Cloud .

    • REGION: האזור שבו יפעל עומס העבודה של עיבוד ברצף (batch processing) ב-Managed Service for Apache Spark.

  2. עוברים אל Lakehouse.

    מעבר אל Lakehouse

  3. בוחרים את הקטלוג LAKEHOUSE_CATALOG_ID ולוחצים על מחיקה.

  4. נכנסים אל Cloud Storage Buckets.

    כניסה לדף Buckets

  5. בוחרים את הדלי ולוחצים על מחיקה.

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