הגדרה של Spark ו-Hive באמצעות קטלוג זמן הריצה של Lakehouse

במאמר הזה מוסבר איך להגדיר את Apache Spark ו-Apache Hive כדי להשתמש בקטלוג של Lakehouse runtime. תלמדו איך ליצור קטלוג של Apache Hive, להגדיר את הסשנים של Spark כדי להתחבר למאגר המטא-נתונים ולהריץ עומסי עבודה כדי ליצור טבלאות שאפשר להריץ עליהן שאילתות ישירות ב-BigQuery.

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

  1. כדי להבין איך Spark מתחבר לקטלוג של Lakehouse runtime, אפשר לקרוא את המאמר About Hive Catalogs in Lakehouse runtime catalog.
  2. פורמטים נתמכים של אחסון וסוגי נתונים
  3. בודקים את המגבלות והשיקולים.
  4. נכנסים לחשבון Google Cloud . אם אתם משתמשים חדשים ב- Google Cloud, צרו חשבון כדי שתוכלו להעריך את הביצועים של המוצרים שלנו בתרחישים מהעולם האמיתי. לקוחות חדשים מקבלים בחינם גם קרדיט בשווי 300$ להרצה, לבדיקה ולפריסה של עומסי העבודה.
  5. Verify that billing is enabled for your Google Cloud project.

  6. Enable the Lakehouse for Apache Iceberg, Dataproc API 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

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

  8. Enable the Lakehouse for Apache Iceberg, Dataproc API 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

התפקידים הנדרשים

כדי לקבל את ההרשאות שדרושות לשימוש בקטלוג של זמן הריצה של Lakehouse, אתם צריכים לבקש מהאדמין להקצות לכם בפרויקט את תפקידי ה-IAM הבאים:

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

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

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

תהליך עבודה כללי

כדי להשתמש בקטלוג של Lakehouse runtime עם Spark ו-Hive, צריך לבצע את תהליך העבודה הכללי הבא:

  1. יצירת Lakehouse for Apache Iceberg לקטלוג Hive.
  2. מגדירים את סשן Spark באמצעות הכלי המועדף (כמו Managed Service for Apache Spark או BigQuery Studio).
  3. ביצוע פעולות על מסדי נתונים וטבלאות בסשן Spark.
  4. שליחת עומסי עבודה של Batch אל Managed Service for Apache Spark ושליחת שאילתות לטבלאות שנוצרות ישירות מ-BigQuery.

יצירת קטלוג Hive של Lakehouse

כדי להשתמש בקטלוג של Lakehouse runtime עם Spark ו-Hive, צריך קודם ליצור קטלוג של Hive.

קטלוג Lakehouse Hive הוא אוסף של מסדי נתונים של Hive. לפני שמריצים משימות Spark, צריך ליצור קטלוג כדי לרשום אותו ב-Lakehouse Metastore. לקטלוג יש שם ומיקום ב-Cloud Storage שבו נמצאים נתוני Hive.

המסוף

  1. במסוף Google Cloud , פותחים את הדף Lakehouse.

    מעבר אל Lakehouse

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

  3. בוחרים באפשרות קטלוג זמן ריצה של Lakehouse.

  4. בשדה Catalog type, בוחרים באפשרות Hive Metastore.

  5. בשדה בחירת קטגוריה של Cloud Storage, מזינים את השם של הקטגוריה של Cloud Storage שבה רוצים להשתמש עם הקטלוג. לחלופין, לוחצים על עיון כדי לבחור באחד ממאגרי הנתונים הקיימים או ליצור מאגר חדש.

  6. בשדה מזהה קטלוג, נותנים שם לקטלוג Lakehouse Hive.

  7. בקטע מיקום ראשי, מציינים את אותו אזור של הקטגוריה.

  8. לוחצים על יצירה.

gcloud

כדי ליצור קטלוג Hive, מריצים את הפקודה הבאה:

gcloud beta biglake hive catalogs create LAKEHOUSE_CATALOG_ID \
    --project=PROJECT_ID \
    --location-uri="gs://GCS_WAREHOUSE_PATH" \
    --primary-location=REGION \
    --description="DESCRIPTION"

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

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

  • GCS_WAREHOUSE_PATH: הנתיב ב-Cloud Storage שבו מאוחסן מחסן הנתונים של Hive.

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

  • REGION: האזור הראשי של ה-Metastore. בקטגוריות באזור יחיד, הוא צריך להיות זהה לאזור של הקטגוריה. בקטגוריות של שני אזורים או של כמה אזורים, צריך לציין אחד מהאזורים שמרכיבים את הקטגוריה, שבו אמורה להיות הרפליקה הראשית של מאגר המידע. האזור השני הופך לעותק המשני.

  • DESCRIPTION: תיאור של הקטלוג.

curl

  1. כדי ליצור קטלוג Hive, מריצים את הפקודה הבאה:
curl -X POST -s -i -H "Authorization: Bearer $(gcloud auth print-access-token)" \
-d '{"locationUri": "gs://GCS_WAREHOUSE_PATH", "description": "DESCRIPTION"}' \
-H "Content-Type:application/json" \ "https://biglake.googleapis.com/hive/v1beta/projects/PROJECT_ID/catalogs?hiveCatalogId=LAKEHOUSE_CATALOG_ID&primary_location=REGION"

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

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

  • GCS_WAREHOUSE_PATH: הנתיב ב-Cloud Storage שבו מאוחסן מחסן הנתונים של Hive.

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

  • REGION: האזור הראשי של ה-Metastore. בקטגוריות באזור יחיד, הוא צריך להיות זהה לאזור של הקטגוריה. בקטגוריות של שני אזורים או של מספר אזורים, צריך לציין אחד מהאזורים שמרכיבים את הקטגוריה, שבו אמורה להיות הרפליקה הראשית של מאגר המטא-נתונים. האזור השני הופך לעותק המשני.

  • DESCRIPTION: תיאור של הקטלוג.

הגדרה ושימוש ב-Spark וב-Hive

כדי להשתמש בקטלוג של סביבת זמן הריצה של Lakehouse, צריך להגדיר את סשן Spark עם מאפיינים ספציפיים. אפשר להגדיר את המאפיינים האלה כשיוצרים אשכול של Managed Service for Apache Spark, או לציין אותם בכל פעם שיוצרים סשן.

המאפיינים האלה כוללים פרטים כמו client factory,‏ Google Cloudמזהה הפרויקט, קטלוג ברירת המחדל וספריית מחסן הנתונים. אחרי שהסשן נוצר, אפשר לבצע פעולות בסיסיות כמו הצגת מסדי נתונים קיימים, יצירת מסדי נתונים חדשים, הגדרת טבלאות והוספת נתונים.

spark-sql

  1. משתמשים ב-SSH כדי להתחבר לצומת הראשי של אשכול Managed Service for Apache Spark.
  2. מריצים את spark-sql בשורת הפקודה עם המאפיינים הבאים כדי להתחיל סשן אינטראקטיבי של Spark SQL:

    spark-sql \
        --conf spark.hive.metastore.client.factory.class=com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory \
        --conf spark.hive.metastore.blms.project.id=PROJECT_ID \
        --conf spark.hive.metastore.blms.catalog.default=LAKEHOUSE_CATALOG_ID \
        --conf spark.hive.metastore.warehouse.dir=gs://GCS_WAREHOUSE_PATH
    
  3. אחרי שהסשן מתחיל, Spark מתחבר לקטלוג של זמן הריצה של Lakehouse.

    מריצים את הפקודות הבאות כדי ליצור משאבים ולשאול עליהם שאילתות:

    
    -- Show all the databases in the current project.
    SHOW DATABASES;
    
    -- Create a database.
    CREATE DATABASE spark_blms_database;
    
    -- Create a Parquet datasource table.
    CREATE TABLE spark_blms_database.parquet_quick_start (id INT, name STRING) USING PARQUET;
    
    -- Insert data into the table.
    INSERT INTO TABLE spark_blms_database.parquet_quick_start VALUES (1, 'my-first-user');
    

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

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

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

    • GCS_WAREHOUSE_PATH: הנתיב ב-Cloud Storage שבו מאוחסן מחסן הנתונים של Hive.

Jupyter Notebook

  1. פועלים לפי ההוראות להפעלת מחברת Jupyter באשכול Managed Service for Apache Spark.
  2. אפשר לגשת לממשק האינטרנט של Jupyter מהכרטיסייה Web interfaces בדף הפרטים של האשכול במסוף Google Cloud .
  3. במחברת חדשה, יוצרים סשן Spark ואז מריצים את השאילתות הבאות:

    from pyspark.sql import SparkSession
    
    # If a Spark session exists, stop it first by running spark.stop()
    spark = SparkSession.builder\
        .master("local")\
        .appName("Lakehouse runtime catalog tutorial")\
        .config("spark.hive.metastore.client.factory.class", "com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory")\
        .config("spark.hive.metastore.blms.project.id", "PROJECT_ID")\
        .config("spark.hive.metastore.warehouse.dir", "gs://GCS_WAREHOUSE_PATH")\
        .config("spark.hive.metastore.blms.catalog.default", "LAKEHOUSE_CATALOG_ID")\
        .getOrCreate()
    
    # Show all the databases.
    df = spark.sql("SHOW DATABASES;")
    df.show()
    
    # Create a database.
    spark.sql("CREATE DATABASE jupyter_blms_db")
    
    # Create a Parquet datasource table.
    spark.sql("CREATE TABLE jupyter_blms_db.parquet_table(id INT, name STRING) USING PARQUET")
    
    # Insert data into the table.
    spark.sql("INSERT INTO TABLE jupyter_blms_db.parquet_table VALUES (1, 'my-first-user');")
    
    # Query from table.
    spark.sql("SELECT * FROM jupyter_blms_db.parquet_table;").show()
    

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

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

    • GCS_WAREHOUSE_PATH: הנתיב ב-Cloud Storage שבו מאוחסן מחסן הנתונים של Hive.

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

מחברת BigQuery

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

    כניסה ל-BigQuery

  2. בחלונית Explorer, לוחצים על + ADD ואז על Python notebook.

  3. בתא קוד, מגדירים את המאפיינים של קטלוג זמן הריצה של Lakehouse ואת הסשן של Managed Service for Apache Spark:

    from google.cloud.dataproc_spark_connect import DataprocSparkSession
    from google.cloud.dataproc_v1 import Session
    import os
    
    os.environ['DATAPROC_SPARK_CONNECT_DEFAULT_DATASOURCE'] = ""
    
    session = Session()
    session.environment_config.execution_config.ttl = {"seconds": 864000}
    session.runtime_config.version = "2.3"
    session.runtime_config.properties = {
     "spark.hive.metastore.blms.project.id": "PROJECT_ID",
     "spark.hive.metastore.blms.catalog.default": "LAKEHOUSE_CATALOG_ID",
     "spark.hive.metastore.warehouse.dir": "gs://GCS_WAREHOUSE_PATH",
     "spark.hive.metastore.client.factory.class": "com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory",
     "spark.sql.catalogImplementation": "hive"
    }
    
    spark = DataprocSparkSession.builder.dataprocSessionConfig(session).getOrCreate()
    print("Spark session created successfully")
    
  4. בתא קוד אחר, יוצרים מסד נתונים וטבלה:

    # Create a database
    spark.sql("CREATE DATABASE bq_spark_blms_database;")
    
    # Create a parquet datasource table
    spark.sql("CREATE TABLE bq_spark_blms_database.parquet_quick_start (id INT, name STRING) USING PARQUET;")
    
    # Insert data into the table
    spark.sql("INSERT INTO TABLE bq_spark_blms_database.parquet_quick_start VALUES (1, 'my-first-user');")
    
    # Query from table
    spark.sql("select * from bq_spark_blms_database.parquet_quick_start;").show()
    

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

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

    • GCS_WAREHOUSE_PATH: הנתיב ב-Cloud Storage שבו מאוחסן מחסן הנתונים של Hive.

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

Hive CLI

כדי להשתמש ב-Hive CLI, צריך להגדיר אותו כך שיתחבר למאגר המטא-נתונים. אפשר לעשות זאת באחת משתי דרכים:

  • אפשרות 1: הגדרה קבועה. מעדכנים את הקובץ /etc/hive/conf/hive-site.xml בצומת הראשי.
  • אפשרות 2: הגדרה זמנית. מספקים את דגלי ההגדרה כשמפעילים את ה-CLI.

כדי להגדיר את CLI ולהריץ את סשן השאילתות, מבצעים את השלבים הבאים:

  1. משתמשים ב-SSH כדי להתחבר לצומת הראשי של אשכול Managed Service for Apache Spark.
  2. בהתאם לשיטת ההגדרה, מבצעים אחת מהפעולות הבאות:

    • כדי להגדיר באופן קבוע (אפשרות 1):

      1. פותחים את /etc/hive/conf/hive-site.xml בצומת הראשי ומוסיפים את המאפיינים הבאים:

        <property>
            <name>hive.metastore.blms.project.id</name>
            <value>PROJECT_ID</value>
            <description></description>
        </property>
        
        <property>
            <name>hive.metastore.client.factory.class</name>
            <value>com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory</value>
            <description></description>
        </property>
        <property>
            <name>hive.metastore.blms.catalog.default</name>
            <value>LAKEHOUSE_CATALOG_ID</value>
            <description></description>
        </property>
        <property>
            <name>hive.metastore.warehouse.dir</name>
            <value>gs://GCS_WAREHOUSE_PATH</value>
            <description></description>
        </property>
        
      2. מפעילים את ה-CLI של Hive:

        hive
        
    • כדי להגדיר באופן זמני בהפעלה (אפשרות 2): מפעילים את Hive CLI עם דגלי הגדרה:

      hive \
      --hiveconf hive.metastore.blms.project.id="PROJECT_ID" \
      --hiveconf hive.metastore.blms.catalog.default="LAKEHOUSE_CATALOG_ID" \
      --hiveconf hive.metastore.client.factory.class=com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory \
      --hiveconf hive.metastore.warehouse.dir="gs://GCS_WAREHOUSE_PATH"
      
  3. לאחר מכן, מריצים את הפקודות הבאות בסשן של Hive CLI:

    show databases;
    create database hive_query_test;
    create table hive_query_test.parquet_table (id INT, name STRING) stored as PARQUET;
    select * from hive_query_test.parquet_table;
    

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

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

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

    • GCS_WAREHOUSE_PATH: הנתיב ב-Cloud Storage שבו מאוחסן מחסן הנתונים של Hive.

שליחת משימה באצווה ב-Managed Service for Apache Spark

אפשר לשלוח עומס עבודה של אצווה ב-PySpark אל Managed Service for Apache Spark שמשתמש בקטלוג של Lakehouse runtime.

קטע הקוד הזה של PySpark מאתחל סשן Spark שמוגדר להתחבר לקטלוג של זמן הריצה של Lakehouse. היא מגדירה מאפיינים חיוניים כמו Lakehouse client factory, Google Cloud מזהה הפרויקט, קטלוג ברירת המחדל וספריית מחסן הנתונים. אחרי שמקימים את הסשן, הקוד מראה איך אפשר להשתמש בפקודות Spark SQL כדי להציג רשימה של מסדי נתונים קיימים, ליצור מסד נתונים חדש ולהגדיר טבלה בפורמט Parquet בתוך מסד הנתונים הזה.

  1. יוצרים קובץ Python עם עבודת PySpark:

    from pyspark.sql import SparkSession
    
    spark = (
        SparkSession.builder.appName("Lakehouse runtime catalog tutorial")
        .config(
           "spark.hive.metastore.client.factory.class",
    "com.google.cloud.bigquery.metastore.client.BigLakeMetastoreClientFactory")
        .config("spark.hive.metastore.blms.project.id", "PROJECT_ID")
        .config("spark.hive.metastore.blms.catalog.default", "LAKEHOUSE_CATALOG_ID")
        .config(
            "spark.hive.metastore.warehouse.dir",
            "gs://GCS_WAREHOUSE_PATH",
        )
        .enableHiveSupport()
        .getOrCreate()
    )
    
    # Show all the databases.
    spark.sql("SHOW DATABASES;").show()
    
    # Create a database.
    spark.sql("CREATE DATABASE dp_serverless_test")
    
    # Create a Parquet datasource table.
    spark.sql(
            "CREATE TABLE dp_serverless_test.parquet_table(id INT, name STRING) USING"
            " PARQUET"
        )
    
    # Query from table.
    spark.sql("SELECT * from  dp_serverless_test.parquet_table").show()
    

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

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

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

    • GCS_WAREHOUSE_PATH: הנתיב ב-Cloud Storage שבו מאוחסן מחסן הנתונים של Hive.

  2. שולחים את המשימה באצווה:

    gcloud dataproc batches submit pyspark PYTHON_SCRIPT_FILE \
        --version=2.2 \
        --project=PROJECT_ID \
        --region=REGION \
        --deps-bucket=gs://CLOUD_STORAGE_BUCKET
    

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

    • PYTHON_SCRIPT_FILE: הנתיב לקובץ האפליקציה של PySpark. זה יכול להיות נתיב מקומי או נתיב אובייקט ב-Cloud Storage.
    • PROJECT_ID: מזהה הפרויקט.
    • REGION: האזור שבו תופעל משימה באצווה.
    • CLOUD_STORAGE_BUCKET: השם של קטגוריית Cloud Storage שמשמשת להעברת תלות של עומסי עבודה.

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

אחרי שיוצרים משאבים מ-Spark בקטלוג של Lakehouse runtime, אפשר להריץ עליהם שאילתות מ-BigQuery Studio.

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

    כניסה ל-BigQuery

  2. מזינים את ההצהרה הבאה בעורך השאילתות:

    SELECT * FROM `PROJECT_ID.LAKEHOUSE_CATALOG_ID.DATABASE_NAME.TABLE_NAME`;
    

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

    • PROJECT_ID: מזהה הפרויקט.
    • LAKEHOUSE_CATALOG_ID: השם של קטלוג Hive.
    • DATABASE_NAME: שם מסד הנתונים.
    • TABLE_NAME: שם הטבלה.