1. סקירה כללית
ב-Lab הזה נסביר איך להגדיר ולהשתמש ב-Apache Spark וב-Jupyter notebooks ב-Managed Apache Spark.
מחברות Jupyter נמצאות בשימוש נרחב לניתוח נתונים לצורך גילוי תובנות ולבניית מודלים של למידת מכונה, כי הן מאפשרות להריץ את הקוד באופן אינטראקטיבי ולראות את התוצאות באופן מיידי.
עם זאת, ההגדרה והשימוש ב-Apache Spark וב-Jupyter Notebooks יכולים להיות מסובכים.

שירות מנוהל ל-Apache Spark מאפשר ליצור אשכול מנוהל של Apache Spark עם Apache Spark, רכיב Jupyter ושער רכיבים תוך כ-90 שניות, כך שהתהליך מהיר וקל.
מה תלמדו
ב-codelab הזה תלמדו איך:
- יצירת קטגוריה של Google Cloud Storage עבור האשכול
- יצירת אשכול מנוהל של Apache Spark עם Jupyter ו-Component Gateway,
- גישה לממשק המשתמש של JupyterLab באינטרנט ב-Managed Apache Spark
- יצירת מחברת באמצעות המחבר של Spark BigQuery Storage
- הפעלת משימת Spark ושרטוט התוצאות.
העלות הכוללת להרצת שיעור ה-Lab הזה ב-Google Cloud היא בערך 1$. פרטים מלאים על התמחור של Managed Apache Spark זמינים כאן.
2. יצירת פרויקט
נכנסים אל Google Cloud Platform Console בכתובת console.cloud.google.com ויוצרים פרויקט חדש:



בשלב הבא, כדי להשתמש במשאבים של Google Cloud, צריך להפעיל את החיוב במסוף Cloud.
העלות של ה-Codelab הזה לא אמורה להיות גבוהה מכמה דולרים, אבל היא יכולה להיות גבוהה יותר אם תחליטו להשתמש ביותר משאבים או אם תשאירו אותם פועלים. בקטע האחרון של שיעור ה-Codelab הזה תלמדו איך לנקות את הפרויקט.
משתמשים חדשים ב-Google Cloud Platform זכאים לתקופת ניסיון בחינם בשווי 300$.
3. הגדרת הסביבה
קודם כל, פותחים את Cloud Shell בלחיצה על הלחצן בפינה השמאלית העליונה של מסוף הענן:

אחרי ש-Cloud Shell נטען, מריצים את הפקודה הבאה כדי להגדיר את מזהה הפרויקט מהשלב הקודם**:**
gcloud config set project <project_id>
אפשר גם לראות את מזהה הפרויקט כשלוחצים על הפרויקט בפינה הימנית העליונה של מסוף הענן:


לאחר מכן, מפעילים את ממשקי ה-API של Managed Apache Spark, Compute Engine ו-BigQuery Storage.
gcloud services enable dataproc.googleapis.com \
compute.googleapis.com \
storage-component.googleapis.com \
bigquery.googleapis.com \
bigquerystorage.googleapis.com
אפשר גם לעשות את זה ב-Cloud Console. לוחצים על סמל התפריט בפינה השמאלית העליונה.

בתפריט הנפתח, בוחרים באפשרות 'API Manager'.

לוחצים על Enable APIs and Services.

מחפשים את ממשקי ה-API הבאים ומפעילים אותם:
- Compute Engine API
- Managed Apache Spark API
- BigQuery API
- BigQuery Storage API
4. יצירת קטגוריה ב-GCS
יוצרים קטגוריה ב-Google Cloud Storage באזור הכי קרוב לנתונים שלכם ונותנים לה שם ייחודי.
הוא ישמש לאשכול המנוהל של Apache Spark.
REGION=us-central1
BUCKET_NAME=<your-bucket-name>
gsutil mb -c standard -l ${REGION} gs://${BUCKET_NAME}
הפלט הבא אמור להתקבל:
Creating gs://<your-bucket-name>/...
5. יצירת אשכול מנוהל של Apache Spark באמצעות Jupyter ו-Component Gateway
יצירת האשכול
הגדרת משתני הסביבה של האשכול
REGION=us-central1
ZONE=us-central1-a
CLUSTER_NAME=spark-jupyter
BUCKET_NAME=<your-bucket-name>
לאחר מכן מריצים את פקודת gcloud הזו כדי ליצור את האשכול עם כל הרכיבים הנדרשים לעבודה עם Jupyter באשכול.
gcloud beta dataproc clusters create ${CLUSTER_NAME} \
--region=${REGION} \
--image-version=2.2 \
--master-machine-type=n1-standard-4 \
--worker-machine-type=n1-standard-4 \
--bucket=${BUCKET_NAME} \
--optional-components=JUPYTER \
--enable-component-gateway
בזמן יצירת האשכול, אמור להופיע הפלט הבא
Waiting on operation [projects/spark-jupyter/regions/us-central1/operations/abcd123456].
Waiting for cluster creation operation...
יצירת האשכול תימשך כ-90 שניות, וכשהוא יהיה מוכן תוכלו לגשת אליו מממשק המשתמש של מסוף Managed Apache Spark Cloud.
בזמן ההמתנה, אפשר להמשיך לקרוא בהמשך כדי לקבל מידע נוסף על הדגלים שמשמשים בפקודת gcloud.
אחרי יצירת האשכול, הפלט הבא אמור להתקבל:
Created [https://dataproc.googleapis.com/v1beta2/projects/project-id/regions/us-central1/clusters/spark-jupyter] Cluster placed in zone [us-central1-a].
הדגלים שמשמשים בפקודה gcloud dataproc create
פירוט של האפשרויות שבהן נעשה שימוש בפקודה gcloud dataproc create
--region=${REGION}
מציינים את האזור והתחום שבהם ייווצר האשכול. כאן אפשר לראות את רשימת האזורים שבהם המינוי זמין.
--image-version=1.4
גרסת התמונה שבה רוצים להשתמש באשכול. כאן אפשר לראות את רשימת הגרסאות הזמינות.
--bucket=${BUCKET_NAME}
מציינים את הקטגוריה של Google Cloud Storage שיצרתם קודם לשימוש באשכול. אם לא תספקו דלי GCS, המערכת תיצור אותו בשבילכם.
כאן גם יישמרו מחברות ה-Notebook שלכם, גם אם תמחקו את האשכול, כי מאגר ה-GCS לא יימחק.
--master-machine-type=n1-standard-4
--worker-machine-type=n1-standard-4
סוגי המכונות שבהן יש להשתמש באשכול Apache Spark המנוהל. כאן אפשר לראות רשימה של סוגי מכונות זמינים.
כברירת מחדל, נוצר צומת ראשי אחד ו-2 צמתים של עובדים אם לא מגדירים את הדגל –num-workers
--optional-components=ANACONDA,JUPYTER
הגדרת הערכים האלה עבור רכיבים אופציונליים תגרום להתקנה של כל הספריות הנדרשות ל-Jupyter ול-Anaconda (שנדרשת למחברות Jupyter) באשכול.
--enable-component-gateway
הפעלת Component Gateway יוצרת קישור ל-App Engine באמצעות Apache Knox ו-Inverting Proxy, שמאפשר גישה קלה, מאובטחת ומאומתת לממשקי האינטרנט של Jupyter ו-JupyterLab. המשמעות היא שכבר לא צריך ליצור מנהרות SSH.
הוא גם ייצור קישורים לכלים אחרים באשכול, כולל Yarn Resource manager ו-Spark History Server, שמועילים להצגת הביצועים של העבודות ודפוסי השימוש באשכול.
6. יצירת מחברת Apache Spark
גישה לממשק האינטרנט של JupyterLab
אחרי שהאשכול מוכן, אפשר למצוא את הקישור של Component Gateway לממשק האינטרנט של JupyterLab. לשם כך, עוברים אל Managed Apache Spark Clusters - Cloud console (אשכולות מנוהלים של Apache Spark – מסוף Cloud), לוחצים על האשכול שיצרתם ועוברים לכרטיסייה Web Interfaces (ממשקי אינטרנט).

תראו שיש לכם גישה ל-Jupyter, שהוא ממשק המחברת הקלאסי, או ל-JupyterLab, שמתואר כממשק המשתמש מהדור הבא של פרויקט Jupyter.
יש הרבה תכונות חדשות ונהדרות בממשק המשתמש של JupyterLab, ולכן אם אתם חדשים בשימוש במסמכי notebook או מחפשים את השיפורים האחרונים, מומלץ להשתמש ב-JupyterLab, כי בסופו של דבר הוא יחליף את הממשק הקלאסי של Jupyter, לפי המסמכים הרשמיים.
יצירת נוטבוק עם ליבת Python 3

בכרטיסייה של מרכז האפליקציות, לוחצים על סמל המחברת של Python 3 כדי ליצור מחברת עם ליבת Python 3 (לא ליבת PySpark). כך אפשר להגדיר את SparkSession במחברת ולכלול את spark-bigquery-connector שנדרש לשימוש ב-BigQuery Storage API.
שינוי השם של הנוטבוק

לוחצים לחיצה ימנית על שם המחברת בסרגל הצד בצד ימין או בסרגל הניווט העליון ומשנים את שם המחברת ל-BigQuery Storage & Spark DataFrames.ipynb.
הרצת קוד Spark ב-notebook

במחברת הזו תשתמשו ב-spark-bigquery-connector, כלי לקריאה ולכתיבה של נתונים בין BigQuery ל-Spark באמצעות BigQuery Storage API.
BigQuery Storage API מביא שיפורים משמעותיים לגישה לנתונים ב-BigQuery באמצעות פרוטוקול מבוסס-RPC. הוא תומך בקריאה ובכתיבה של נתונים במקביל, וגם בפורמטים שונים של סריאליזציה כמו Apache Avro ו-Apache Arrow. ברמה גבוהה, המשמעות היא שיפור משמעותי בביצועים, במיוחד כשמדובר במערכי נתונים גדולים יותר.
בתא הראשון, בודקים את גרסת Scala של האשכול כדי לכלול את הגרסה הנכונה של קובץ ה-jar של מחבר spark-bigquery.
קלט [1]:
!scala -version
Output [1]:
יצירת סשן Spark וצירוף חבילת המחבר spark-bigquery.
אם גרסת Scala שלכם היא 2.11, צריך להשתמש בחבילה הבאה.
com.google.cloud.spark:spark-bigquery-with-dependencies_2.11:0.15.1-beta
אם גרסת Scala שלכם היא 2.12, צריך להשתמש בחבילה הבאה.
com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.15.1-beta
קלט [2]:
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName('BigQuery Storage & Spark DataFrames') \
.config('spark.jars.packages', 'com.google.cloud.spark:spark-bigquery-with-dependencies_2.11:0.15.1-beta') \
.getOrCreate()
הפעלת repl.eagerEval
הפעולה הזו תציג את התוצאות של DataFrames בכל שלב בלי הצורך החדש להציג df.show(), וגם תשפר את הפורמט של הפלט.
קלט [3]:
spark.conf.set("spark.sql.repl.eagerEval.enabled",True)
קריאת טבלה של BigQuery לתוך Spark DataFrame
יוצרים Spark DataFrame על ידי קריאת נתונים ממערך נתונים ציבורי של BigQuery. התהליך הזה משתמש ב- מחבר spark-bigquery וב-BigQuery Storage API כדי לטעון את הנתונים אל אשכול Spark.
יוצרים Spark DataFrame וטוענים נתונים ממערך הנתונים הציבורי של BigQuery לגבי צפיות בדפים בוויקיפדיה. שימו לב שלא מריצים שאילתה על הנתונים, כי משתמשים ב-מחבר spark-bigquery כדי לטעון את הנתונים ל-Spark, שם יתבצע עיבוד הנתונים. כשמריצים את הקוד הזה, הטבלה לא נטענת בפועל כי מדובר בהערכה עצלה ב-Spark, והביצוע יתרחש בשלב הבא.
קלט [4]:
table = "bigquery-public-data.wikipedia.pageviews_2020"
df_wiki_pageviews = spark.read \
.format("bigquery") \
.option("table", table) \
.option("filter", "datehour >= '2020-03-01' AND datehour < '2020-03-02'") \
.load()
df_wiki_pageviews.printSchema()
פלט [4]:

בוחרים את העמודות הנדרשות ומחילים מסנן באמצעות where() שהוא כינוי ל-filter().
כשמריצים את הקוד הזה, מופעלת פעולת Spark והנתונים נקראים מ-BigQuery Storage בשלב הזה.
קלט [5]:
df_wiki_en = df_wiki_pageviews \
.select("datehour", "wiki", "views") \
.where("views > 1000 AND wiki in ('en', 'en.m')") \
df_wiki_en
פלט [5]:

כדי לראות את הדפים המובילים, מקבצים לפי כותרת וממיינים לפי צפיות בדף
קלט [6]:
import pyspark.sql.functions as F
df_datehour_totals = df_wiki_en \
.groupBy("datehour") \
.agg(F.sum('views').alias('total_views'))
df_datehour_totals.orderBy('total_views', ascending=False)
פלט [6]:
7. שימוש בספריות של Python לשרטוט ב-notebook
אתם יכולים להשתמש בספריות שונות של Python לשרטוט כדי לשרטט את הפלט של משימות Spark.
המרת Spark DataFrame ל-Pandas DataFrame
ממירים את Spark DataFrame ל-Pandas DataFrame ומגדירים את datehour כאינדקס. האפשרות הזו שימושית אם רוצים לעבוד עם הנתונים ישירות ב-Python ולשרטט את הנתונים באמצעות ספריות השרטוט הרבות שזמינות ב-Python.
קלט [7]:
spark.conf.set("spark.sql.execution.arrow.enabled", "true")
pandas_datehour_totals = df_datehour_totals.toPandas()
pandas_datehour_totals.set_index('datehour', inplace=True)
pandas_datehour_totals.head()
פלט [7]:

Plotting Pandas Dataframe
מייבאים את ספריית matplotlib שנדרשת להצגת התרשימים במחברת
קלט [8]:
import matplotlib.pyplot as plt
משתמשים בפונקציית התרשים של Pandas כדי ליצור תרשים קו מ-Pandas DataFrame.
קלט [9]:
pandas_datehour_totals.plot(kind='line',figsize=(12,6));
פלט [9]:
בדיקה שה-notebook נשמר ב-GCS
עכשיו אמור להיות לכם מסמך Jupyter notebook ראשון שפועל באשכול Managed Apache Spark. נותנים שם למחברת, והיא תינשמר אוטומטית בקטגוריית GCS שבה השתמשתם כשיצרתם את האשכול.
אפשר לבדוק את זה באמצעות פקודת gsutil הבאה ב-Cloud Shell
BUCKET_NAME=<your-bucket-name>
gsutil ls gs://${BUCKET_NAME}/notebooks/jupyter
הפלט הבא אמור להתקבל:
gs://bucket-name/notebooks/jupyter/
gs://bucket-name/notebooks/jupyter/BigQuery Storage & Spark DataFrames.ipynb
8. טיפ לאופטימיזציה – שמירת נתונים במטמון בזיכרון
יכול להיות שיהיו תרחישים שבהם תרצו שהנתונים יהיו בזיכרון במקום לקרוא אותם מ-BigQuery Storage בכל פעם.
המשימה הזו תקרא את הנתונים מ-BigQuery ותדחוף את המסנן ל-BigQuery. לאחר מכן, הצבירה תחושב ב-Apache Spark.
import pyspark.sql.functions as F
table = "bigquery-public-data.wikipedia.pageviews_2020"
df_wiki_pageviews = spark.read \
.format("bigquery") \
.option("table", table) \
.option("filter", "datehour >= '2020-03-01' AND datehour < '2020-03-02'") \
.load()
df_wiki_en = df_wiki_pageviews \
.select("title", "wiki", "views") \
.where("views > 10 AND wiki in ('en', 'en.m')")
df_wiki_en_totals = df_wiki_en \
.groupBy("title") \
.agg(F.sum('views').alias('total_views'))
df_wiki_en_totals.orderBy('total_views', ascending=False)
אפשר לשנות את העבודה שלמעלה כך שתכלול מטמון של הטבלה, ועכשיו המסנן בעמודה של הוויקי יוחל בזיכרון על ידי Apache Spark.
import pyspark.sql.functions as F
table = "bigquery-public-data.wikipedia.pageviews_2020"
df_wiki_pageviews = spark.read \
.format("bigquery") \
.option("table", table) \
.option("filter", "datehour >= '2020-03-01' AND datehour < '2020-03-02'") \
.load()
df_wiki_all = df_wiki_pageviews \
.select("title", "wiki", "views") \
.where("views > 10")
# cache the data in memory
df_wiki_all.cache()
df_wiki_en = df_wiki_all \
.where("wiki in ('en', 'en.m')")
df_wiki_en_totals = df_wiki_en \
.groupBy("title") \
.agg(F.sum('views').alias('total_views'))
df_wiki_en_totals.orderBy('total_views', ascending=False)
אחר כך תוכלו לסנן לפי שפה אחרת של ויקיפדיה באמצעות הנתונים שבמטמון, במקום לקרוא שוב נתונים מאחסון BigQuery, ולכן הפעולה תהיה מהירה יותר.
df_wiki_de = df_wiki_all \
.where("wiki in ('de', 'de.m')")
df_wiki_de_totals = df_wiki_de \
.groupBy("title") \
.agg(F.sum('views').alias('total_views'))
df_wiki_de_totals.orderBy('total_views', ascending=False)
אפשר להסיר את המטמון באמצעות הפקודה
df_wiki_all.unpersist()
9. דוגמאות ל-notebook לתרחישי שימוש נוספים
ב-Managed Apache Spark GitHub repo יש מחברות Jupyter עם דפוסי Apache Spark נפוצים לטעינת נתונים, לשמירת נתונים ולשרטוט הנתונים באמצעות מגוון מוצרים של Google Cloud Platform וכלים בקוד פתוח:
10. הסרת המשאבים
כדי להימנע מחיובים מיותרים בחשבון GCP אחרי שתסיימו את המדריך למתחילים הזה:
- מוחקים את הקטגוריה של Cloud Storage עבור הסביבה שיצרתם
- מחיקת סביבת Apache Spark מנוהלת
אם יצרתם פרויקט רק בשביל ה-Codelab הזה, אתם יכולים גם למחוק אותו:
- במסוף GCP, נכנסים לדף Projects.
- ברשימת הפרויקטים, בוחרים את הפרויקט שרוצים למחוק ולוחצים על מחיקה.
- בתיבה, כותבים את מזהה הפרויקט ולוחצים על Shut down כדי למחוק את הפרויקט.
רישיון
עבודה זו מורשית תחת רישיון Creative Commons שמותנה בייחוס 3.0 כללי, ורישיון Apache 2.0.