‫Apache Spark ו-Jupyter Notebooks ב-Managed Service for Apache Spark

1. סקירה כללית

ב-Lab הזה נסביר איך להגדיר ולהשתמש ב-Apache Spark וב-Jupyter notebooks ב-Managed Apache Spark.

מחברות Jupyter נמצאות בשימוש נרחב לניתוח נתונים לצורך גילוי תובנות ולבניית מודלים של למידת מכונה, כי הן מאפשרות להריץ את הקוד באופן אינטראקטיבי ולראות את התוצאות באופן מיידי.

עם זאת, ההגדרה והשימוש ב-Apache Spark וב-Jupyter Notebooks יכולים להיות מסובכים.

b9ed855863c57d6.png

שירות מנוהל ל-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 ויוצרים פרויקט חדש:

7e541d932b20c074.png2deefc9295d114ea.pnga92a49afe05008a.png

בשלב הבא, כדי להשתמש במשאבים של Google Cloud, צריך להפעיל את החיוב במסוף Cloud.

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

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

3. הגדרת הסביבה

קודם כל, פותחים את Cloud Shell בלחיצה על הלחצן בפינה השמאלית העליונה של מסוף הענן:

a10c47ee6ca41c54.png

אחרי ש-Cloud Shell נטען, מריצים את הפקודה הבאה כדי להגדיר את מזהה הפרויקט מהשלב הקודם**:**

gcloud config set project <project_id>

אפשר גם לראות את מזהה הפרויקט כשלוחצים על הפרויקט בפינה הימנית העליונה של מסוף הענן:

b4b233632ce0c3c4.pngc7e39ffc6dec3765.png

לאחר מכן, מפעילים את ממשקי ה-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. לוחצים על סמל התפריט בפינה השמאלית העליונה.

2bfc27ef9ba2ec7d.png

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

408af5f32c4b7c25.png

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

a9c0e84296a7ba5b.png

מחפשים את ממשקי ה-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 (ממשקי אינטרנט).

afc40202d555de47.png

תראו שיש לכם גישה ל-Jupyter, שהוא ממשק המחברת הקלאסי, או ל-JupyterLab, שמתואר כממשק המשתמש מהדור הבא של פרויקט Jupyter.

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

יצירת נוטבוק עם ליבת Python 3

a463623f2ebf0518.png

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

שינוי השם של הנוטבוק

196a3276ed07e1f3.png

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

הרצת קוד Spark ב-notebook

fbac38062e5bb9cf.png

במחברת הזו תשתמשו ב-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]:f580e442576b8b1f.png יצירת סשן 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]:

c107a33f6fc30ca.png

בוחרים את העמודות הנדרשות ומחילים מסנן באמצעות 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]:

ad363cbe510d625a.png

כדי לראות את הדפים המובילים, מקבצים לפי כותרת וממיינים לפי צפיות בדף

קלט [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]:f718abd05afc0f4.png

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]:

3df2aaa2351f028d.png

Plotting Pandas Dataframe

מייבאים את ספריית matplotlib שנדרשת להצגת התרשימים במחברת

קלט [8]:

import matplotlib.pyplot as plt

משתמשים בפונקציית התרשים של Pandas כדי ליצור תרשים קו מ-Pandas DataFrame.

קלט [9]:

pandas_datehour_totals.plot(kind='line',figsize=(12,6));

פלט [9]:bade7042c3033594.png

בדיקה שה-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 אחרי שתסיימו את המדריך למתחילים הזה:

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

אם יצרתם פרויקט רק בשביל ה-Codelab הזה, אתם יכולים גם למחוק אותו:

  1. במסוף GCP, נכנסים לדף Projects.
  2. ברשימת הפרויקטים, בוחרים את הפרויקט שרוצים למחוק ולוחצים על מחיקה.
  3. בתיבה, כותבים את מזהה הפרויקט ולוחצים על Shut down כדי למחוק את הפרויקט.

רישיון

עבודה זו מורשית תחת רישיון Creative Commons שמותנה בייחוס 3.0 כללי, ורישיון Apache 2.0.