אפשר להשתמש ב-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 CLI:
חישוב העלויות
במדריך הזה נעשה שימוש ברכיבים של Google Cloudשחלים עליהם חיובים, כולל:
- Managed Service for Apache Spark
- BigQuery
- Cloud Storage
אפשר להשתמש במחשבון עלויות כדי ליצור הערכת עלויות בהתאם לשימוש החזוי.
הגדרת חיוב
כברירת מחדל, הפרויקט שמשויך לפרטי הכניסה או לחשבון השירות מחויב על השימוש ב-API. כדי לחייב פרויקט אחר, מגדירים את מאפיין התצורה הבא: spark.conf.set("parentProject", "<BILLED-GCP-PROJECT>").
אפשר גם להוסיף את המאפיין הזה לפעולת קריאה או כתיבה, באופן הבא:
.option("parentProject", "<BILLED-GCP-PROJECT>").
שליחת עומס עבודה של ספירת מילים ב-PySpark
בדוגמה הזו, הנתונים נקראים מ-BigQuery לתוך Spark DataFrame כדי לבצע ספירת מילים באמצעות API של מקור נתונים סטנדרטי.
המחבר כותב את הפלט של ספירת המילים ל-BigQuery ברצף הפעולות הבא:
מאחסנת את הנתונים בקבצים זמניים בקטגוריה של Cloud Storage
העתקת הנתונים בפעולה אחת מהקטגוריה של Cloud Storage אל BigQuery
הקובץ הזמני נמחק ב-Cloud Storage אחרי השלמת פעולת הטעינה ב-BigQuery (הקובץ הזמני נמחק גם אחרי סיום האפליקציה של Spark). אם המחיקה נכשלת, צריך למחוק קבצים זמניים לא רצויים ב-Cloud Storage, שבדרך כלל נמצאים בתיקייה
gs://BUCKET_NAME/.spark-bigquery-JOB_ID-UUID.
השלבים להפעלת עומס העבודה של ספירת המילים
- פותחים טרמינל מקומי או Cloud Shell.
- יוצרים את
wordcount_datasetבאמצעות כלי שורת הפקודה bq בטרמינל מקומי או ב-Cloud Shell.bq mk wordcount_dataset
- יוצרים קטגוריה של Cloud Storage באמצעות Google Cloud CLI.
מחליפים אתgcloud storage buckets create gs://BUCKET_NAME
BUCKET_NAMEבשם של קטגוריה של Cloud Storage שיצרתם. - יוצרים את הקובץ
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')
- שליחת עומס העבודה של קבוצת 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ולוחצים על תצוגה מקדימה.
איור 1: תצוגה מקדימה של טבלת הפלט ב-BigQuery
למידע נוסף
- BigQuery Storage & Spark SQL - Python
- יצירת קובץ הגדרת טבלה למקור נתונים חיצוני
- שימוש בנתונים שחולקו למחיצות באופן חיצוני