שימוש במחבר Spark BigQuery

אפשר להשתמש ב-spark-bigquery-connector עם Managed Service for Apache Spark כדי לקרוא ולכתוב נתונים מ-BigQuery ואליו. במדריך הזה מוצגת אפליקציית PySpark שמשתמשת ב-spark-bigquery-connector.

אישור גרסת המחבר

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

הפיכת המחבר לזמין לאפליקציות (אם צריך)

מחבר BigQuery מותקן בכל גרסאות זמן הריצה הנתמכות של Managed Service for Apache Spark. אם אתם משתמשים בגרסת זמן ריצה לא נתמכת שלא מתקינה את המחבר (Spark runtime 1.0), אתם יכולים להפוך את המחבר לזמין לאפליקציה באחת משתי הדרכים הבאות:

  • משתמשים בפרמטר jars כדי להפנות לקובץ jar של מחבר כששולחים עומס עבודה של אצווה ב-Managed Service for Apache Spark או כשמריצים סשן אינטראקטיבי. בדוגמה הבאה של עומס עבודה באצווה מצוין קובץ JAR של מחבר (ברשימת קובצי ה-JAR של המחברים הזמינים אפשר לעיין במאגר GoogleCloudDataproc/spark-bigquery-connector ב-GitHub).
    • דוגמה ל-Google Cloud CLI:
      gcloud dataproc batches submit pyspark \
          --region=REGION \
          --jars=spark-3.5-bigquery-version.jar \
          ... other args
      

חישוב העלויות

במדריך הזה נעשה שימוש ברכיבים של Google Cloudשחלים עליהם חיובים, כולל:

  • Managed Service for Apache Spark
  • BigQuery
  • Cloud Storage

אפשר להשתמש במחשבון עלויות כדי ליצור הערכת עלויות בהתאם לשימוש החזוי.

משתמשים חדשים ב-Cloud Platform עשויים להיות זכאים לתקופת ניסיון בחינם.

הגדרת חיוב

כברירת מחדל, הפרויקט שמשויך לפרטי הכניסה או לחשבון השירות מחויב על השימוש ב-API. כדי לחייב פרויקט אחר, מגדירים את מאפיין התצורה הבא: spark.conf.set("parentProject", "<BILLED-GCP-PROJECT>").

אפשר גם להוסיף את המאפיין הזה לפעולת קריאה או כתיבה, באופן הבא: .option("parentProject", "<BILLED-GCP-PROJECT>").

שליחת עומס עבודה של ספירת מילים ב-PySpark

בדוגמה הזו, הנתונים נקראים מ-BigQuery לתוך Spark DataFrame כדי לבצע ספירת מילים באמצעות API של מקור נתונים סטנדרטי.

המחבר כותב את הפלט של ספירת המילים ל-BigQuery ברצף הפעולות הבא:

  1. מאחסנת את הנתונים בקבצים זמניים בקטגוריה של Cloud Storage

  2. העתקת הנתונים בפעולה אחת מהקטגוריה של Cloud Storage אל BigQuery

  3. הקובץ הזמני נמחק ב-Cloud Storage אחרי השלמת פעולת הטעינה ב-BigQuery (הקובץ הזמני נמחק גם אחרי סיום האפליקציה של Spark). אם המחיקה נכשלת, צריך למחוק קבצים זמניים לא רצויים ב-Cloud Storage, שבדרך כלל נמצאים בתיקייה gs://BUCKET_NAME/.spark-bigquery-JOB_ID-UUID.

השלבים להפעלת עומס העבודה של ספירת המילים

  1. פותחים טרמינל מקומי או Cloud Shell.
  2. יוצרים את wordcount_dataset באמצעות כלי שורת הפקודה bq בטרמינל מקומי או ב-Cloud Shell.
    bq mk wordcount_dataset
    
  3. יוצרים קטגוריה של Cloud Storage באמצעות Google Cloud CLI.
    gcloud storage buckets create gs://BUCKET_NAME
    
    מחליפים את BUCKET_NAME בשם של קטגוריה של Cloud Storage שיצרתם.
  4. יוצרים את הקובץ wordcount.py באופן מקומי בעורך טקסט על ידי העתקת קוד PySpark הבא.
    #!/usr/bin/python
    """BigQuery I/O PySpark example."""
    from pyspark.sql import SparkSession
    
    spark = SparkSession \
      .builder \
      .appName('spark-bigquery-demo') \
      .getOrCreate()
    
    # Cloud Storage bucket used by the connector for temporary BigQuery
    # export data.
    bucket = "BUCKET_NAME"
    spark.conf.set('temporaryGcsBucket', bucket)
    
    # Load data from BigQuery.
    words = spark.read.format('bigquery') \
      .load('bigquery-public-data.samples.shakespeare') \
      .load()
    words.createOrReplaceTempView('words')
    
    # Perform word count.
    word_count = spark.sql(
        'SELECT word, SUM(word_count) AS word_count FROM words GROUP BY word')
    word_count.show()
    word_count.printSchema()
    
    # Save the data to BigQuery
    word_count.write.format('bigquery') \
      .save('wordcount_dataset.wordcount_output')
  5. שליחת עומס העבודה של קבוצת PySpark:
    gcloud dataproc batches submit pyspark wordcount.py \
        --region=REGION \
        --deps-bucket=BUCKET_NAME
    
    פלט לדוגמה של מסוף:
    ...
    +---------+----------+
    |     word|word_count|
    +---------+----------+
    |     XVII|         2|
    |    spoil|        28|
    |    Drink|         7|
    |forgetful|         5|
    |   Cannot|        46|
    |    cures|        10|
    |   harder|        13|
    |  tresses|         3|
    |      few|        62|
    |  steel'd|         5|
    | tripping|         7|
    |   travel|        35|
    |   ransom|        55|
    |     hope|       366|
    |       By|       816|
    |     some|      1169|
    |    those|       508|
    |    still|       567|
    |      art|       893|
    |    feign|        10|
    +---------+----------+
    only showing top 20 rows
    
    root
     |-- word: string (nullable = false)
     |-- word_count: long (nullable = true)
    

    כדי לראות תצוגה מקדימה של טבלת הפלט במסוף Google Cloud , פותחים את הדף BigQuery, בוחרים את הטבלה wordcount_output ולוחצים על תצוגה מקדימה.
    עיבוד של תצוגה מקדימה של טבלה ב-BigQuery
    איור 1: תצוגה מקדימה של טבלת הפלט ב-BigQuery

למידע נוסף