יצירת צינור עיבוד נתונים של Dataflow באמצעות Python
במאמר הזה מוסבר איך להשתמש ב-Apache Beam SDK for Python כדי ליצור תוכנית שמגדירה צינור עיבוד נתונים. לאחר מכן מריצים את צינור עיבוד הנתונים באמצעות רץ מקומי ישיר או רץ מבוסס-ענן כמו Dataflow. בסרטון איך משתמשים ב-WordCount ב-Apache Beam יש מבוא לצינור WordCount.
לחצו על תראו לי איך כדי לקרוא הסבר מפורט על המשימה ישירות במסוף Google Cloud :
לפני שמתחילים
- נכנסים לחשבון Google Cloud . אם אתם משתמשים חדשים ב- Google Cloud, צרו חשבון כדי שתוכלו להעריך את הביצועים של המוצרים שלנו בתרחישים מהעולם האמיתי. לקוחות חדשים מקבלים בחינם גם קרדיט בשווי 300$ להרצה, לבדיקה ולפריסה של עומסי העבודה.
-
התקינו את ה-CLI של Google Cloud.
-
אם אתם משתמשים בספק זהויות חיצוני (IdP), קודם אתם צריכים להיכנס ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.
-
כדי לאתחל את ה-CLI של gcloud, הריצו את הפקודה הבאה:
gcloud init -
יוצרים או בוחרים Google Cloud פרויקט.
תפקידים שנדרשים כדי לבחור או ליצור פרויקט
- Select a project: כדי לבחור פרויקט לא צריך תפקיד IAM ספציפי – אפשר לבחור כל פרויקט שקיבלתם בו תפקיד.
-
יצירת פרויקט: כדי ליצור פרויקט, צריך את התפקיד Project Creator (
roles/resourcemanager.projectCreator), שכולל את ההרשאהresourcemanager.projects.create. איך מקצים תפקידים
-
יוצרים Google Cloud פרויקט:
gcloud projects create PROJECT_ID
מחליפים את
PROJECT_IDבשם של פרויקט Google Cloud שיוצרים. -
בוחרים את הפרויקט שיצרתם: Google Cloud
gcloud config set project PROJECT_ID
מחליפים את
PROJECT_IDבשם הפרויקט ב- Google Cloud .
מפעילים את ממשקי ה-API של Dataflow, Compute Engine, Cloud Logging, Cloud Storage, Google Cloud Storage JSON, BigQuery, Cloud Pub/Sub, Cloud Datastore ו-Cloud Resource Manager, אם הם עדיין לא מופעלים:
תפקידים שנדרשים להפעלת ממשקי API
כדי להפעיל ממשקי API, נדרשת ההרשאה
serviceusage.services.enable. אם יצרתם את הפרויקט, סביר להניח שכבר יש לכם את ההרשאה הזו דרך התפקיד 'בעלים' (roles/owner). אחרת, תוכלו לקבל את ההרשאה הזו דרך התפקיד 'אדמין בממשק Service Usage' (roles/serviceusage.serviceUsageAdmin). איך מקצים תפקידיםgcloud services enable dataflow
compute_component logging storage_component storage_api bigquery pubsub datastore.googleapis.com cloudresourcemanager.googleapis.com -
יוצרים פרטי כניסה לאימות מקומי עבור חשבון המשתמש:
gcloud auth application-default login
אם מוחזרת שגיאת אימות ואתם משתמשים בספק זהויות חיצוני (IdP), ודאו ש נכנסתם ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.
-
מעניקים תפקידים לחשבון המשתמש. מריצים את הפקודה הבאה לכל אחד מהתפקידים הבאים ב-IAM:
roles/iam.supportUser, roles/datastream.admin, roles/monitoring.metricsScopesViewer, roles/cloudaicompanion.settingsAdmingcloud projects add-iam-policy-binding PROJECT_ID --member="user:USER_IDENTIFIER" --role=ROLE
מחליפים את מה שכתוב בשדות הבאים:
-
PROJECT_ID: מזהה הפרויקט. -
USER_IDENTIFIER: המזהה של חשבון המשתמש . לדוגמה,myemail@example.com. -
ROLE: תפקיד ה-IAM שאתם מקצים לחשבון המשתמש.
-
-
התקינו את ה-CLI של Google Cloud.
-
אם אתם משתמשים בספק זהויות חיצוני (IdP), קודם אתם צריכים להיכנס ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.
-
כדי לאתחל את ה-CLI של gcloud, הריצו את הפקודה הבאה:
gcloud init -
יוצרים או בוחרים Google Cloud פרויקט.
תפקידים שנדרשים כדי לבחור או ליצור פרויקט
- Select a project: כדי לבחור פרויקט לא צריך תפקיד IAM ספציפי – אפשר לבחור כל פרויקט שקיבלתם בו תפקיד.
-
יצירת פרויקט: כדי ליצור פרויקט, צריך את התפקיד Project Creator (
roles/resourcemanager.projectCreator), שכולל את ההרשאהresourcemanager.projects.create. איך מקצים תפקידים
-
יוצרים Google Cloud פרויקט:
gcloud projects create PROJECT_ID
מחליפים את
PROJECT_IDבשם של פרויקט Google Cloud שיוצרים. -
בוחרים את הפרויקט שיצרתם: Google Cloud
gcloud config set project PROJECT_ID
מחליפים את
PROJECT_IDבשם הפרויקט ב- Google Cloud .
מפעילים את ממשקי ה-API של Dataflow, Compute Engine, Cloud Logging, Cloud Storage, Google Cloud Storage JSON, BigQuery, Cloud Pub/Sub, Cloud Datastore ו-Cloud Resource Manager, אם הם עדיין לא מופעלים:
תפקידים שנדרשים להפעלת ממשקי API
כדי להפעיל ממשקי API, נדרשת ההרשאה
serviceusage.services.enable. אם יצרתם את הפרויקט, סביר להניח שכבר יש לכם את ההרשאה הזו דרך התפקיד 'בעלים' (roles/owner). אחרת, תוכלו לקבל את ההרשאה הזו דרך התפקיד 'אדמין בממשק Service Usage' (roles/serviceusage.serviceUsageAdmin). איך מקצים תפקידיםgcloud services enable dataflow
compute_component logging storage_component storage_api bigquery pubsub datastore.googleapis.com cloudresourcemanager.googleapis.com -
יוצרים פרטי כניסה לאימות מקומי עבור חשבון המשתמש:
gcloud auth application-default login
אם מוחזרת שגיאת אימות ואתם משתמשים בספק זהויות חיצוני (IdP), ודאו ש נכנסתם ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.
-
מעניקים תפקידים לחשבון המשתמש. מריצים את הפקודה הבאה לכל אחד מהתפקידים הבאים ב-IAM:
roles/iam.supportUser, roles/datastream.admin, roles/monitoring.metricsScopesViewer, roles/cloudaicompanion.settingsAdmingcloud projects add-iam-policy-binding PROJECT_ID --member="user:USER_IDENTIFIER" --role=ROLE
מחליפים את מה שכתוב בשדות הבאים:
-
PROJECT_ID: מזהה הפרויקט. -
USER_IDENTIFIER: המזהה של חשבון המשתמש . לדוגמה,myemail@example.com. -
ROLE: תפקיד ה-IAM שאתם מקצים לחשבון המשתמש.
-
מקצים תפקידים לחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine. מריצים את הפקודה הבאה לכל אחד מהתפקידים הבאים ב-IAM:
roles/dataflow.adminroles/dataflow.workerroles/storage.objectAdmin
gcloud projects add-iam-policy-binding PROJECT_ID --member="serviceAccount:PROJECT_NUMBER-compute@developer.gserviceaccount.com" --role=SERVICE_ACCOUNT_ROLE
- מחליפים את
PROJECT_IDבמזהה הפרויקט. - מחליפים את
PROJECT_NUMBERבמספר הפרויקט. כדי למצוא את מספר הפרויקט, אפשר לעיין במאמר בנושא זיהוי פרויקטים או להשתמש בפקודהgcloud projects describe. - מחליפים את
SERVICE_ACCOUNT_ROLEבכל אחד מהתפקידים.
-
יוצרים קטגוריה של Cloud Storage ומגדירים אותה כך:
-
מגדירים את סוג האחסון (storage class) לאפשרות הבאה:
S(Standard). -
מגדירים את מיקום האחסון לאזור הבא:
US(ארצות הברית). -
מחליפים את
BUCKET_NAMEבשם ייחודי לקטגוריה. שם הקטגוריה לא יכול להכיל מידע רגיש כי מרחב השמות של הקטגוריות זמין וגלוי לכולם.
gcloud storage buckets create gs://BUCKET_NAME --default-storage-class STANDARD --location US
-
מגדירים את סוג האחסון (storage class) לאפשרות הבאה:
- מעתיקים את Google Cloud מזהה הפרויקט ואת שם הקטגוריה של Cloud Storage. תצטרכו את הערכים האלה בהמשך המאמר.
מגדירים את הסביבה
בקטע הזה, משתמשים בשורת הפקודה כדי להגדיר סביבה וירטואלית מבודדת של Python להרצת פרויקט צינור הנתונים באמצעות venv. התהליך הזה מאפשר לבודד את יחסי התלות של פרויקט אחד מיחסי התלות של פרויקטים אחרים.
אם אין לכם שורת פקודה זמינה, אתם יכולים להשתמש ב-Cloud Shell. מנהל החבילות של Python 3 כבר מותקן ב-Cloud Shell, כך שאפשר לדלג לשלב של יצירת סביבה וירטואלית.
כדי להתקין את Python ואז ליצור סביבה וירטואלית, פועלים לפי השלבים הבאים:
- בודקים ש-Python 3 ו-
pipפועלים במערכת:python --version python -m pip --version
- אם נדרש, מתקינים את Python 3 ואז מגדירים סביבה וירטואלית של Python: פועלים לפי ההוראות שבקטעים התקנת Python והגדרת venv בדף הגדרת סביבת פיתוח של Python.
אחרי שמסיימים את המדריך למתחילים, אפשר להשבית את הסביבה הווירטואלית על ידי הפעלת הפקודה deactivate.
הורדה של Apache Beam SDK
Apache Beam SDK הוא מודל תכנות בקוד פתוח לצינורות נתונים. מגדירים צינור עיבוד נתונים באמצעות תוכנית Apache Beam ואז בוחרים רץ, כמו Dataflow, כדי להריץ את צינור עיבוד הנתונים.
כדי להוריד ולהתקין את Apache Beam SDK:
- מוודאים שאתם בסביבה הווירטואלית של Python שיצרתם בקטע הקודם.
מוודאים שההנחיה מתחילה ב-
<env_name>, כאשרenv_nameהוא השם של הסביבה הווירטואלית. - מתקינים את הגרסה האחרונה של Apache Beam SDK ל-Python:
pip install apache-beam[gcp]
הפעלת הפייפליין באופן מקומי
כדי לראות איך צינור פועל באופן מקומי, משתמשים במודול Python מוכן מראש לדוגמה wordcount שכלולה בחבילה apache_beam.
בדוגמה לצינור העיבוד wordcount מתבצעות הפעולות הבאות:
מקבל קובץ טקסט כקלט.
קובץ הטקסט הזה נמצא בקטגוריה של Cloud Storage עם שם המשאב
gs://dataflow-samples/shakespeare/kinglear.txt.- מנתח כל שורה למילים.
- מבצעת ספירת תדירות של המילים שעברו טוקניזציה.
כדי להכין את צינור הנתונים wordcount באופן מקומי:
- בטרמינל המקומי, מריצים את הדוגמה
wordcount:python -m apache_beam.examples.wordcount \ --output outputs
- צופים בפלט של צינור עיבוד הנתונים:
more outputs* - כדי לצאת, מקישים על q.
wordcount.pyקוד המקור
ב-Apache Beam GitHub.
הרצת הפייפליין בשירות Dataflow
בקטע הזה, מריצים את צינור העיבוד לדוגמהwordcount מחבילת apache_beam בשירות Dataflow. בדוגמה הזו, הערך DataflowRunner מוגדר כפרמטר של --runner.
- מריצים את הפייפליין:
python -m apache_beam.examples.wordcount \ --region DATAFLOW_REGION \ --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://BUCKET_NAME/results/outputs \ --runner DataflowRunner \ --project PROJECT_ID \ --temp_location gs://BUCKET_NAME/tmp/
מחליפים את מה שכתוב בשדות הבאים:
-
DATAFLOW_REGION: האזור שבו רוצים לפרוס את עבודת Dataflow, לדוגמה:europe-west1הדגל
--regionמבטל את אזור ברירת המחדל שמוגדר בשרת המטא-נתונים, בלקוח המקומי או במשתני הסביבה. -
BUCKET_NAME: שם הקטגוריה ב-Cloud Storage שהעתקתם קודם -
PROJECT_ID: Google Cloud מזהה הפרויקט שהעתקתם קודם
-
צפייה בתוצאות
כשמריצים צינור באמצעות Dataflow, התוצאות מאוחסנות בקטגוריה של Cloud Storage. בקטע הזה, מוודאים שהצינור פועל באמצעות מסוף Google Cloud או מסוף מקומי.
מסוףGoogle Cloud
כדי לראות את התוצאות ב Google Cloud מסוף, פועלים לפי השלבים הבאים:
- נכנסים לדף Jobs ב-Dataflow במסוף Google Cloud .
בדף Jobs מוצגים פרטים על
wordcountהמשימה, כולל הסטטוס שלה. בהתחלה הסטטוס יהיה Running ואחר כך Succeeded. - נכנסים לדף Buckets של Cloud Storage.
ברשימת הקטגוריות בפרויקט, לוחצים על קטגוריית האחסון שיצרתם קודם.
בספרייה
wordcountמוצגים קובצי הפלט שנוצרו על ידי העבודה.
מסוף מקומי
אפשר לראות את התוצאות בטרמינל או באמצעות Cloud Shell.
- כדי לראות את רשימת קובצי הפלט, משתמשים בפקודה
gcloud storage ls:gcloud storage ls gs://BUCKET_NAME/results/outputs* --long
- כדי לראות את התוצאות בקובצי הפלט, משתמשים בפקודה
gcloud storage cat:gcloud storage cat gs://BUCKET_NAME/results/outputs*
מחליפים את BUCKET_NAME בשם של קטגוריית Cloud Storage שבה נעשה שימוש בתוכנית של צינור העיבוד.
שינוי הקוד של צינור עיבוד הנתונים
בצינור העיבודwordcount בדוגמאות הקודמות יש הבחנה בין מילים באותיות רישיות לבין מילים באותיות קטנות.
בשלבים הבאים מוסבר איך לשנות את צינור הנתונים כך שצינור הנתונים wordcount לא יהיה תלוי באותיות רישיות.
- במחשב המקומי, מורידים את העותק העדכני של הקוד
wordcountממאגר Apache Beam ב-GitHub. - בטרמינל המקומי, מריצים את צינור העיבוד:
python wordcount.py --output outputs
- מעיינים בתוצאות:
more outputs* - כדי לצאת, מקישים על q.
- פותחים את הקובץ
wordcount.pyבכלי עריכה לבחירתכם. - בתוך הפונקציה
run, בודקים את השלבים בפייפליין:counts = ( lines | 'Split' >> (beam.ParDo(WordExtractingDoFn()).with_output_types(str)) | 'PairWithOne' >> beam.Map(lambda x: (x, 1)) | 'GroupAndSum' >> beam.CombinePerKey(sum))
אחרי
split, השורות מפוצלות למילים כמחרוזות. - כדי להפוך את המחרוזות לאותיות קטנות, משנים את השורה אחרי
split: השינוי הזה ממפה את הפונקציהcounts = ( lines | 'Split' >> (beam.ParDo(WordExtractingDoFn()).with_output_types(str)) | 'lowercase' >> beam.Map(str.lower) | 'PairWithOne' >> beam.Map(lambda x: (x, 1)) | 'GroupAndSum' >> beam.CombinePerKey(sum))
str.lowerלכל מילה. השורה הזו שווה ל-beam.Map(lambda word: str.lower(word)). - שומרים את הקובץ ומריצים את משימת
wordcountששיניתם:python wordcount.py --output outputs
- צפייה בתוצאות של הפייפליין ששונה:
more outputs* - כדי לצאת, מקישים על q.
- מריצים את הפייפליין ששיניתם בשירות Dataflow:
python wordcount.py \ --region DATAFLOW_REGION \ --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://BUCKET_NAME/results/outputs \ --runner DataflowRunner \ --project PROJECT_ID \ --temp_location gs://BUCKET_NAME/tmp/
מחליפים את מה שכתוב בשדות הבאים:
-
DATAFLOW_REGION: האזור שבו רוצים לפרוס את משימת Dataflow -
BUCKET_NAME: שם הקטגוריה שלכם ב-Cloud Storage -
PROJECT_ID: מזהה הפרויקט ב- Google Cloud
-
הסרת המשאבים
כדי לא לצבור חיובים בחשבון Google Cloud על המשאבים שבהם השתמשתם בדף הזה, אתם צריכים למחוק את הפרויקט Google Cloud יחד עם המשאבים.
- במסוף Google Cloud , נכנסים לדף Buckets של Cloud Storage.
- לוחצים על תיבת הסימון של הקטגוריה שרוצים למחוק.
- כדי למחוק את הקטגוריה, לוחצים על Delete ופועלים לפי ההוראות.
אם אתם משאירים את הפרויקט, אתם צריכים לבטל את התפקידים שהקציתם לחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine. מריצים את הפקודה הבאה לכל אחד מהתפקידים הבאים ב-IAM:
roles/dataflow.adminroles/dataflow.workerroles/storage.objectAdmin
gcloud projects remove-iam-policy-binding PROJECT_ID \ --member=serviceAccount:PROJECT_NUMBER-compute@developer.gserviceaccount.com \ --role=SERVICE_ACCOUNT_ROLE
-
אם תרצו, תוכלו לבטל את פרטי הכניסה שיצרתם ולמחוק את הקובץ המקומי של פרטי הכניסה.
gcloud auth application-default revoke
-
אם רוצים, מבטלים את פרטי הכניסה של ה-CLI של gcloud.
gcloud auth revoke