Managed Service for Apache Spark

1. סקירה כללית – Serverless for Apache Spark

Managed Service for Apache Spark הוא שירות מנוהל במלואו וניתן להתאמה לעומס, להרצת Apache Spark, ‏ Apache Flink, ‏ Presto ועוד הרבה כלים ומסגרות בקוד פתוח. שימוש ב-Managed Apache Spark לעדכון אגם הנתונים (data lake), ל-ETL / ELT ולמדעי נתונים מאובטחים, בקנה מידה גלובלי. ‫Apache Spark מנוהל משולב באופן מלא גם עם כמה שירותים של Google Cloud, כולל BigQuery, ‏ Cloud Storage, ‏ Gemini Enterprise Agent Engine ו-Knowledge Catalog.

‫Apache Spark מנוהל זמין בשתי גרסאות:

  • ‫Apache Spark מנוהל ללא שרת מאפשר להריץ משימות PySpark בלי צורך להגדיר תשתית והתאמה אוטומטית לעומס. ‫Apache Spark מנוהל תומך בעומסי עבודה של PySpark batch ובסשנים או ב-Notebooks.
  • אשכולות מנוהלים של Apache Spark מאפשרים לכם לנהל אשכול Hadoop YARN לעומסי עבודה של Spark שמבוססים על YARN, בנוסף לכלים בקוד פתוח כמו Flink ו-Presto. אתם יכולים להתאים את האשכולות מבוססי הענן שלכם עם קנה מידה אנכי או אופקי ככל שתרצו, כולל שינוי קנה מידה אוטומטי.

בשיעור Codelab הזה תלמדו על כמה דרכים שונות שבהן אפשר להשתמש ב-Dataproc Serverless.

‫Apache Spark נבנה במקור להפעלה באשכולות Hadoop, והשתמש ב-YARN כמנהל המשאבים שלו. תחזוקה של אשכולות Hadoop דורשת מומחיות ספציפית, וצריך לוודא שהרבה הגדרות שונות באשכולות מוגדרות בצורה נכונה. בנוסף, יש עוד קבוצה נפרדת של הגדרות שמשתמשים צריכים להגדיר ב-Spark. כתוצאה מכך, יש הרבה תרחישים שבהם מפתחים משקיעים יותר זמן בהגדרת התשתית שלהם במקום לעבוד על קוד Spark עצמו.

‫Dataproc Serverless מייתר את הצורך בהגדרה ידנית של אשכולות Hadoop או Spark. ‫Dataproc Serverless לא פועל ב-Hadoop ומשתמש בהקצאת משאבים דינמית משלו כדי לקבוע את דרישות המשאבים שלו, כולל התאמה אוטומטית לעומס. עדיין אפשר להתאים אישית קבוצת משנה קטנה של מאפייני Spark באמצעות Dataproc Serverless, אבל ברוב המקרים לא תצטרכו לבצע שינויים כאלה.

2. הגדרה

תתחילו בהגדרת הסביבה והמשאבים שבהם נעשה שימוש ב-codelab הזה.

יוצרים פרויקט ב-Google Cloud. אפשר להשתמש באחד קיים.

פותחים את Cloud Shell בלחיצה על הסמל שלו בסרגל הכלים של מסוף Cloud.

ba0bb17945a73543.png

‫Cloud Shell מספק סביבת מעטפת מוכנה לשימוש שאפשר להשתמש בה בשיעור Codelab הזה.

68c4ebd2a8539764.png

שם הפרויקט מוגדר ב-Cloud Shell כברירת מחדל. כדי לוודא זאת, מריצים את הפקודה echo $GOOGLE_CLOUD_PROJECT. אם מזהה הפרויקט לא מופיע בפלט, צריך להגדיר אותו.

export GOOGLE_CLOUD_PROJECT=<your-project-id>

הגדרת אזור ב-Compute Engine למשאבים, כמו us-central1 או europe-west2.

export REGION=<your-region>

הפעלת ממשקי ה-API

בשיעור Codelab הזה נעשה שימוש בממשקי ה-API הבאים:

  • BigQuery
  • Dataproc

מפעילים את ממשקי ה-API הנדרשים. הפעולה תימשך כדקה, וכשתסתיים תוצג הודעה על הצלחה.

gcloud services enable bigquery.googleapis.com
gcloud services enable dataproc.googleapis.com

הגדרת גישה לרשת

כדי להשתמש ב-Dataproc Serverless, צריך להפעיל את הגישה הפרטית ל-Google באזור שבו מריצים את משימות Spark, כי לדרייברים ולמבצעים של Spark יש רק כתובות IP פרטיות. מריצים את הפקודה הבאה כדי להפעיל אותה בתת-הרשת default.

gcloud compute networks subnets update default \
  --region=${REGION} \
  --enable-private-ip-google-access

כדי לוודא ש-Google Private Access מופעל, מריצים את הפקודה הבאה. הפלט יהיה True או False.

gcloud compute networks subnets describe default \
  --region=${REGION} \
  --format="get(privateIpGoogleAccess)"

יצירה של קטגוריית אחסון

יוצרים קטגוריית אחסון שתשמש לאחסון נכסים שנוצרו ב-codelab הזה.

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

export BUCKET=<your-bucket-name>

יוצרים את הקטגוריה באזור שבו רוצים להריץ את משימות Spark.

gsutil mb -l ${REGION} gs://${BUCKET}

אפשר לראות שהקטגוריה שלכם זמינה במסוף Cloud Storage. אפשר גם להריץ את הפקודה gsutil ls כדי לראות את הדלי.

יצירת שרת היסטוריה מתמשך

ממשק המשתמש של Spark מספק מערך עשיר של כלי ניפוי באגים ותובנות לגבי משימות Spark. כדי לראות את ממשק המשתמש של Spark למשימות שהושלמו ב-Dataproc Serverless, צריך ליצור אשכול Dataproc עם צומת יחיד כדי להשתמש בו כשרת היסטוריה מתמשך.

מגדירים שם לשרת ההיסטוריה הקבוע.

PHS_CLUSTER_NAME=my-phs

מריצים את הפקודה הבאה.

gcloud dataproc clusters create ${PHS_CLUSTER_NAME} \
    --region=${REGION} \
    --single-node \
    --enable-component-gateway \
    --properties=spark:spark.history.fs.logDirectory=gs://${BUCKET}/phs/*/spark-job-history

בהמשך ה-codelab נסביר יותר על ממשק המשתמש של Spark ועל שרת ההיסטוריה המתמשך.

3. הרצת משימות Serverless Spark באמצעות Dataproc Batches

בדוגמה הזו תעבדו עם קבוצת נתונים מתוך מערך הנתונים הציבורי של נסיעות באופניים של Citi בניו יורק. ‫NYC Citi Bikes היא מערכת שיתוף אופניים בתשלום בניו יורק. תבצעו כמה טרנספורמציות פשוטות ותדפיסו את עשרת מזהי התחנות הכי פופולריים של Citi Bike. בדוגמה הזו נעשה שימוש גם ב-spark-bigquery-connector, מחבר קוד פתוח, כדי לקרוא ולכתוב נתונים בצורה חלקה בין Spark ל-BigQuery.

משכפלים את מאגר Github הבא, cd, לתוך הספרייה שמכילה את הקובץ citibike.py.

git clone https://github.com/GoogleCloudPlatform/devrel-demos.git
cd devrel-demos/data-analytics/next-2022-workshop/dataproc-serverless

citibike.py

import sys

from pyspark.sql import SparkSession
from pyspark.sql.functions import col
from pyspark.sql.types import BooleanType

if len(sys.argv) == 1:
    print("Please provide a GCS bucket name.")

bucket = sys.argv[1]
table = "bigquery-public-data:new_york_citibike.citibike_trips"

spark = SparkSession.builder \
          .appName("pyspark-example") \
          .config("spark.jars","gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar") \
          .getOrCreate()

df = spark.read.format("bigquery").load(table)

top_ten = df.filter(col("start_station_id") \
            .isNotNull()) \
            .groupBy("start_station_id") \
            .count() \
            .orderBy("count", ascending=False) \
            .limit(10) \
            .cache()

top_ten.show()

top_ten.write.option("header", True).csv(f"gs://{bucket}/citibikes_top_ten_start_station_ids")

שולחים את העבודה ל-Serverless Spark באמצעות Cloud SDK, שזמין כברירת מחדל ב-Cloud Shell. מריצים את הפקודה הבאה במעטפת, שמשתמשת ב-Cloud SDK וב-Dataproc Batches API כדי לשלוח משימות Serverless Spark.

gcloud dataproc batches submit pyspark citibike.py \
  --batch=citibike-job \
  --region=${REGION} \
  --deps-bucket=gs://${BUCKET} \
  --jars=gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar \
--history-server-cluster=projects/${GOOGLE_CLOUD_PROJECT}/regions/${REGION}/clusters/${PHS_CLUSTER_NAME} \
  -- ${BUCKET}

כדי להסביר את זה:

  • gcloud dataproc batches submit מתייחס אל Dataproc Batches API.
  • pyspark מציין שאתם שולחים משימת PySpark.
  • --batch הוא שם המשימה. אם לא תציינו מזהה, המערכת תשתמש במזהה UUID שנוצר באופן אקראי.
  • --region=${REGION} הוא האזור הגיאוגרפי שבו העבודה תעובד.
  • --deps-bucket=${BUCKET} הוא המקום שאליו קובץ ה-Python המקומי שלכם מועלה לפני שהוא מופעל בסביבה בלי שרת (serverless).
  • --jars=gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar כולל את קובץ ה-jar של spark-bigquery-connector בסביבת זמן הריצה של Spark.
  • --history-server-cluster=projects/${GOOGLE_CLOUD_PROJECT}/regions/${REGION}/clusters/${PHS_CLUSTER} הוא השם המוגדר במלואו של שרת ההיסטוריה הקבוע. כאן מאוחסנים נתוני האירועים של Spark (בנפרד מהפלט של המסוף), וניתן לראות אותם בממשק המשתמש של Spark.
  • התווים -- בסוף השורה מציינים שכל מה שמעבר לזה יהיה ארגומנטים של זמן הריצה של התוכנית. במקרה הזה, שולחים את שם הקטגוריה, כפי שנדרש בעבודה.

הפלט הבא יוצג כשמגישים את הקבוצה.

Batch [citibike-job] submitted.

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

+----------------+------+
|start_station_id| count|
+----------------+------+
|             519|551078|
|             497|423334|
|             435|403795|
|             426|384116|
|             293|372255|
|             402|367194|
|             285|344546|
|             490|330378|
|             151|318700|
|             477|311403|
+----------------+------+

Batch [citibike-job] finished.

בקטע הבא נסביר איך לאתר את היומנים של העבודה הזו.

תכונות נוספות

עם Spark Serverless, יש לכם אפשרויות נוספות להרצת העבודות.

  • אתם יכולים ליצור תמונת Docker מותאמת אישית שהעבודה שלכם תפעל עליה. זו דרך מצוינת לכלול יחסי תלות נוספים, כולל ספריות Python ו-R.
  • אתם יכולים לקשר מופע של Dataproc Metastore לעבודה כדי לגשת למטא-נתונים של Hive.
  • כדי לקבל שליטה נוספת, Dataproc Serverless תומך בהגדרה של קבוצה קטנה של מאפייני Spark.

4. מדדים וניראות ב-Dataproc

ב-Dataproc Batches Console מופיעים כל העבודות של Dataproc Serverless. במסוף אפשר לראות את מזהה האצווה, המיקום והסטטוס של כל עבודה, את זמן היצירה והזמן שחלף ואת הסוג. כדי לראות מידע נוסף על עבודה מסוימת, לוחצים על מזהה האצווה שלה.

בדף הזה מוצג מידע כמו Monitoring (מעקב), שבו אפשר לראות כמה Batch Spark Executors (מנהלי ביצוע של Spark באצווה) נעשה שימוש בעבודה לאורך זמן (מה שמצביע על מידת ההתאמה האוטומטית של קנה המידה).

בכרטיסייה פרטים מוצגים עוד מטא-נתונים על העבודה, כולל ארגומנטים ופרמטרים שנשלחו עם העבודה.

אפשר לגשת לכל היומנים גם מהדף הזה. כשמריצים משימות Dataproc Serverless, נוצרים שלושה סוגים שונים של יומנים:

  • ברמת השירות
  • פלט המסוף
  • רישום ביומן של אירועים ב-Spark

ברמת השירות, כולל יומנים שנוצרו על ידי שירות Dataproc Serverless. הדוגמאות כוללות בקשות של Dataproc בלי שרת (serverless) למעבדים (CPU) נוספים לצורך התאמה אוטומטית לעומס (automatic scaling). כדי לראות את היומנים האלה, לוחצים על הצגת היומנים. היומנים ייפתחו ב-Cloud Logging.

אפשר לראות את הפלט של המסוף בקטע פלט.זהו הפלט שנוצר על ידי העבודה, כולל מטא-נתונים ש-Spark מדפיס כשהעבודה מתחילה או כל הצהרת הדפסה שמשולבת בעבודה.

אפשר לגשת לרישום האירועים ב-Spark דרך ממשק המשתמש של Spark. מכיוון שסיפקתם ל-Spark job שרת היסטוריה קבוע, אתם יכולים לגשת לממשק המשתמש של Spark בלחיצה על View Spark History Server (הצגת שרת ההיסטוריה של Spark), שכולל מידע על Spark jobs שהופעלו בעבר. מידע נוסף על ממשק המשתמש של Spark זמין במסמכי התיעוד הרשמיים של Spark.

5. תבניות Dataproc: BQ -> GCS

Dataproc Templates הם כלים בקוד פתוח שעוזרים לפשט עוד יותר את משימות עיבוד הנתונים בענן. הם משמשים כעטיפה ל-Dataproc Serverless וכוללים תבניות להרבה משימות ייבוא וייצוא נתונים, כולל:

  • BigQuerytoGCS וגם GCStoBigQuery
  • GCStoBigTable
  • GCStoJDBC וגם JDBCtoGCS
  • HivetoBigQuery
  • MongotoGCS וגם GCStoMongo

הרשימה המלאה זמינה בקובץ README.

בקטע הזה נשתמש ב-Dataproc Templates כדי לייצא נתונים מ-BigQuery ל-GCS.

שכפול המאגר

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

git clone https://github.com/GoogleCloudPlatform/dataproc-templates.git
cd dataproc-templates/python

הגדרת הסביבה

עכשיו מגדירים משתני סביבה. תבניות Dataproc משתמשות במשתנה הסביבה GCP_PROJECT למזהה הפרויקט, לכן צריך להגדיר אותו כך: GOOGLE_CLOUD_PROJECT.

export GCP_PROJECT=${GOOGLE_CLOUD_PROJECT}

האזור צריך להיות מוגדר בסביבה מהשלבים הקודמים. אם לא, צריך להגדיר אותו כאן.

export REGION=<region>

תבניות Dataproc משתמשות ב-spark-bigquery-conector לעיבוד משימות BigQuery, וצריך לכלול את ה-URI במשתנה סביבה JARS. מגדירים את המשתנה JARS.

export JARS="gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.26.0.jar"

הגדרת פרמטרים של תבנית

מגדירים את השם של מאגר זמני לשימוש השירות.

export GCS_STAGING_LOCATION=gs://${BUCKET}

לאחר מכן, מגדירים כמה משתנים ספציפיים למשימה. לגבי טבלת הקלט, שוב תהיה הפניה למערך הנתונים NYC Citibike ב-BigQuery.

BIGQUERY_GCS_INPUT_TABLE=bigquery-public-data.new_york_citibike.citibike_trips

אפשר לבחור מבין האפשרויות הבאות: csv, ‏parquet, ‏avro או json. ב-Codelab הזה, בוחרים באפשרות CSV. בקטע הבא מוסבר איך להשתמש ב-Dataproc Templates כדי להמיר סוגי קבצים.

BIGQUERY_GCS_OUTPUT_FORMAT=csv

מגדירים את מצב הפלט לoverwrite. אפשר לבחור בין overwrite, append, ignore או errorifexists.

BIGQUERY_GCS_OUTPUT_MODE=overwrite

מגדירים את מיקום הפלט ב-GCS כנתיב בקטגוריה.

BIGQUERY_GCS_OUTPUT_LOCATION=gs://${BUCKET}/BQtoGCS

הרצת התבנית

מריצים את תבנית BIGQUERYTOGCS על ידי ציון שלה למטה והזנת פרמטרי הקלט שהגדרתם.

./bin/start.sh \
-- --template=BIGQUERYTOGCS \
        --bigquery.gcs.input.table=${BIGQUERY_GCS_INPUT_TABLE} \
        --bigquery.gcs.output.format=${BIGQUERY_GCS_OUTPUT_FORMAT} \
        --bigquery.gcs.output.mode=${BIGQUERY_GCS_OUTPUT_MODE} \
        --bigquery.gcs.output.location=${BIGQUERY_GCS_OUTPUT_LOCATION}

הפלט יהיה די רועש, אבל אחרי דקה בערך תראו את הדברים הבאים.

Batch [5766411d6c78444cb5e80f305308d8f8] submitted.
...
Batch [5766411d6c78444cb5e80f305308d8f8] finished.

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

gsutil ls ${BIGQUERY_GCS_OUTPUT_LOCATION}

כברירת מחדל, Spark כותב לכמה קבצים, בהתאם לכמות הנתונים. במקרה כזה, ייווצרו בערך 30 קבצים. שמות קובצי הפלט של Spark מעוצבים עם part- ואחריו מספר בן חמש ספרות (שמציין את מספר החלק) ומחרוזת גיבוב. בדרך כלל, כשמדובר בכמויות גדולות של נתונים, Spark כותב לכמה קבצים. לדוגמה, שם הקובץ part-00000-cbf69737-867d-41cc-8a33-6521a725f7a0-c000.csv.

6. Dataproc Templates: CSV to parquet

מעכשיו תשתמשו ב-Dataproc Templates כדי להמיר נתונים ב-GCS מסוג קובץ אחד לסוג קובץ אחר באמצעות GCSTOGCS. התבנית הזו משתמשת ב-SparkSQL ומספקת אפשרות לשלוח גם שאילתת SparkSQL לעיבוד במהלך ההמרה, לעיבוד נוסף.

אישור משתני הסביבה

מוודאים שהערכים GCP_PROJECT,‏ REGION ו-GCS_STAGING_BUCKET מוגדרים מהקטע הקודם.

echo ${GCP_PROJECT}
echo ${REGION}
echo ${GCS_STAGING_LOCATION}

הגדרת פרמטרים של תבנית

עכשיו מגדירים פרמטרים של הגדרות ל-GCStoGCS. מתחילים עם המיקום של קובצי הקלט. שימו לב: זוהי ספרייה ולא קובץ ספציפי, כי כל הקבצים בספרייה יעברו עיבוד. מגדירים את הערך BIGQUERY_GCS_OUTPUT_LOCATION.

GCS_TO_GCS_INPUT_LOCATION=${BIGQUERY_GCS_OUTPUT_LOCATION}

מגדירים את הפורמט של קובץ הקלט.

GCS_TO_GCS_INPUT_FORMAT=csv

מגדירים את פורמט הפלט הרצוי. אפשר לבחור בפורמט parquet, ‏ json, ‏ avro או csv.

GCS_TO_GCS_OUTPUT_FORMAT=parquet

מגדירים את מצב הפלט לoverwrite. אפשר לבחור בין overwrite, append, ignore או errorifexists.

GCS_TO_GCS_OUTPUT_MODE=overwrite

מגדירים את מיקום הפלט.

GCS_TO_GCS_OUTPUT_LOCATION=gs://${BUCKET}/GCStoGCS

הרצת התבנית

מריצים את התבנית GCStoGCS.

./bin/start.sh \
-- --template=GCSTOGCS \
        --gcs.to.gcs.input.location=${GCS_TO_GCS_INPUT_LOCATION} \
        --gcs.to.gcs.input.format=${GCS_TO_GCS_INPUT_FORMAT} \
        --gcs.to.gcs.output.format=${GCS_TO_GCS_OUTPUT_FORMAT} \
        --gcs.to.gcs.output.mode=${GCS_TO_GCS_OUTPUT_MODE} \
        --gcs.to.gcs.output.location=${GCS_TO_GCS_OUTPUT_LOCATION}

הפלט יהיה די רועש, אבל אחרי דקה בערך אמורה להופיע הודעת הצלחה כמו זו שבהמשך.

Batch [c198787ba8e94abc87e2a0778c05ec8a] submitted.
...
Batch [c198787ba8e94abc87e2a0778c05ec8a] finished.

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

gsutil ls ${GCS_TO_GCS_OUTPUT_LOCATION}

בעזרת התבנית הזו, אפשר גם לספק שאילתות SparkSQL על ידי העברת gcs.to.gcs.temp.view.name ו-gcs.to.gcs.sql.query לתבנית, וכך להריץ שאילתת SparkSQL על הנתונים לפני הכתיבה ל-GCS.

7. פינוי משאבים

כדי להימנע מחיובים מיותרים בחשבון GCP אחרי שתסיימו את ה-codelab הזה:

  1. מוחקים את הקטגוריה של Cloud Storage עבור הסביבה שיצרתם.
gsutil rm -r gs://${BUCKET}
  1. מוחקים את אשכול Dataproc שמשמש את שרת ההיסטוריה המתמשך.
gcloud dataproc clusters delete ${PHS_CLUSTER_NAME} \
  --region=${REGION}
  1. מוחקים את המשימות של Dataproc Serverless. עוברים אל מסוף אצווה, מסמנים את התיבה לצד כל עבודה שרוצים למחוק ולוחצים על מחיקה.

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

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

8. המאמרים הבאים

במקורות המידע הבאים אתם יכולים לקרוא על דרכים נוספות ליהנות מ-Serverless Spark: