שימוש במאגרי תגים בהתאמה אישית עם ספריות C++‎

במדריך הזה תיצרו צינור עיבוד נתונים שמשתמש במאגרי תמונות מותאמים אישית עם ספריות C++‎ כדי להריץ תהליך עבודה מקביל מאוד של Dataflow HPC. במדריך הזה מוסבר איך להשתמש ב-Dataflow וב-Apache Beam כדי להריץ אפליקציות של מחשוב רשתי שדורשות חלוקת נתונים לפונקציות שפועלות על ליבות רבות.

במדריך הזה נדגים איך להריץ את הפייפליין קודם באמצעות Direct Runner ואחר כך באמצעות Dataflow Runner. הפעלת צינור מקומית מאפשרת לבדוק את הצינור לפני הפריסה שלו.

בדוגמה הזו נעשה שימוש בקישורי Cython ובפונקציות מספריית GMP. לא משנה באיזו ספרייה או כלי לאיגוד אתם משתמשים, אתם יכולים להחיל את אותם העקרונות על צינור הנתונים שלכם.

קוד הדוגמה זמין ב-GitHub.

מטרות

  • יצירת צינור עיבוד נתונים שמשתמש בקונטיינרים מותאמים אישית עם ספריות C++‎.

  • יצירת קובץ אימג' של קונטיינר Docker באמצעות קובץ Docker.

  • אורזים את הקוד ואת יחסי התלות בקונטיינר Docker.

  • מריצים את הפייפליין באופן מקומי כדי לבדוק אותו.

  • הפעלת צינור הנתונים בסביבה מבוזרת.

עלויות

במסמך הזה משתמשים ברכיבים הבאים של Google Cloud, והשימוש בהם כרוך בתשלום:

  • Artifact Registry
  • Cloud Build
  • Cloud Storage
  • Compute Engine
  • Dataflow

כדי להעריך את ההוצאות בהתאם לתחזית השימוש שלכם, אתם יכולים להיעזר במחשבון העלויות.

משתמשים חדשים של Google Cloud ? יכול להיות שאתם זכאים לתקופת ניסיון בחינם.

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

לפני שמתחילים

  1. נכנסים לחשבון Google Cloud . אם אתם משתמשים חדשים ב- Google Cloud, צרו חשבון כדי שתוכלו להעריך את הביצועים של המוצרים שלנו בתרחישים מהעולם האמיתי. לקוחות חדשים מקבלים בחינם גם קרדיט בשווי 300$ להרצה, לבדיקה ולפריסה של עומסי העבודה.
  2. התקינו את ה-CLI של Google Cloud.

  3. אם אתם משתמשים בספק זהויות חיצוני (IdP), קודם אתם צריכים להיכנס ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.

  4. כדי לאתחל את ה-CLI של gcloud, הריצו את הפקודה הבאה:

    gcloud init
  5. יוצרים או בוחרים 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 .

  6. מוודאים שהחיוב מופעל בפרויקט Google Cloud .

  7. מפעילים את ממשקי ה-API של Cloud Storage,‏ Cloud Storage JSON,‏ Compute Engine,‏ Dataflow,‏ Resource Manager,‏ Artifact Registry ו-Cloud Build:

    תפקידים שנדרשים להפעלת ממשקי API

    כדי להפעיל ממשקי API, נדרשת ההרשאה serviceusage.services.enable. אם יצרתם את הפרויקט, סביר להניח שכבר יש לכם את ההרשאה הזו דרך התפקיד 'בעלים' (roles/owner). אחרת, תוכלו לקבל את ההרשאה הזו דרך התפקיד 'אדמין של Service Usage' (roles/serviceusage.serviceUsageAdmin). איך מקצים תפקידים

    gcloud services enable compute.googleapis.com dataflow.googleapis.com storage_component storage_api cloudresourcemanager.googleapis.com artifactregistry.googleapis.com cloudbuild.googleapis.com
  8. יוצרים פרטי כניסה לאימות מקומי עבור חשבון המשתמש:

    gcloud auth application-default login

    אם מוחזרת שגיאת אימות ואתם משתמשים בספק זהויות חיצוני (IdP), ודאו ש נכנסתם ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.

  9. מעניקים תפקידים לחשבון המשתמש. מריצים את הפקודה הבאה לכל אחד מהתפקידים הבאים ב-IAM: roles/iam.serviceAccountUser

    gcloud projects add-iam-policy-binding PROJECT_ID --member="user:USER_IDENTIFIER" --role=ROLE

    מחליפים את מה שכתוב בשדות הבאים:

    • PROJECT_ID: מזהה הפרויקט.
    • USER_IDENTIFIER: המזהה של חשבון המשתמש . לדוגמה, myemail@example.com.
    • ROLE: תפקיד ה-IAM שאתם מקצים לחשבון המשתמש.
  10. התקינו את ה-CLI של Google Cloud.

  11. אם אתם משתמשים בספק זהויות חיצוני (IdP), קודם אתם צריכים להיכנס ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.

  12. כדי לאתחל את ה-CLI של gcloud, הריצו את הפקודה הבאה:

    gcloud init
  13. יוצרים או בוחרים 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 .

  14. מוודאים שהחיוב מופעל בפרויקט Google Cloud .

  15. מפעילים את ממשקי ה-API של Cloud Storage,‏ Cloud Storage JSON,‏ Compute Engine,‏ Dataflow,‏ Resource Manager,‏ Artifact Registry ו-Cloud Build:

    תפקידים שנדרשים להפעלת ממשקי API

    כדי להפעיל ממשקי API, נדרשת ההרשאה serviceusage.services.enable. אם יצרתם את הפרויקט, סביר להניח שכבר יש לכם את ההרשאה הזו דרך התפקיד 'בעלים' (roles/owner). אחרת, תוכלו לקבל את ההרשאה הזו דרך התפקיד 'אדמין של Service Usage' (roles/serviceusage.serviceUsageAdmin). איך מקצים תפקידים

    gcloud services enable compute.googleapis.com dataflow.googleapis.com storage_component storage_api cloudresourcemanager.googleapis.com artifactregistry.googleapis.com cloudbuild.googleapis.com
  16. יוצרים פרטי כניסה לאימות מקומי עבור חשבון המשתמש:

    gcloud auth application-default login

    אם מוחזרת שגיאת אימות ואתם משתמשים בספק זהויות חיצוני (IdP), ודאו ש נכנסתם ל-CLI של gcloud באמצעות המאגר המאוחד לניהול זהויות.

  17. מעניקים תפקידים לחשבון המשתמש. מריצים את הפקודה הבאה לכל אחד מהתפקידים הבאים ב-IAM: roles/iam.serviceAccountUser

    gcloud projects add-iam-policy-binding PROJECT_ID --member="user:USER_IDENTIFIER" --role=ROLE

    מחליפים את מה שכתוב בשדות הבאים:

    • PROJECT_ID: מזהה הפרויקט.
    • USER_IDENTIFIER: המזהה של חשבון המשתמש . לדוגמה, myemail@example.com.
    • ROLE: תפקיד ה-IAM שאתם מקצים לחשבון המשתמש.
  18. יוצרים חשבון שירות של עובד שמנוהל על ידי משתמש עבור צינור עיבוד הנתונים החדש ומקצים לחשבון השירות את התפקידים הנדרשים.

    1. כדי ליצור את חשבון השירות, מריצים את הפקודה gcloud iam service-accounts create:

      gcloud iam service-accounts create parallelpipeline \
          --description="Highly parallel pipeline worker service account" \
          --display-name="Highly parallel data pipeline access"
    2. נותנים לחשבון השירות תפקידים. מריצים את הפקודה הבאה לכל אחד מהתפקידים הבאים ב-IAM:

      • roles/dataflow.admin
      • roles/dataflow.worker
      • roles/storage.objectAdmin
      • roles/artifactregistry.reader
      gcloud projects add-iam-policy-binding PROJECT_ID --member="serviceAccount:parallelpipeline@PROJECT_ID.iam.gserviceaccount.com" --role=SERVICE_ACCOUNT_ROLE

      מחליפים את SERVICE_ACCOUNT_ROLE בכל אחד מהתפקידים.

    3. נותנים לחשבון Google תפקיד שמאפשר ליצור אסימוני גישה לחשבון השירות:

      gcloud iam service-accounts add-iam-policy-binding parallelpipeline@PROJECT_ID.iam.gserviceaccount.com --member="user:EMAIL_ADDRESS" --role=roles/iam.serviceAccountTokenCreator

הורדת הקוד לדוגמה ושינוי הספריות

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

משכפלים את מאגר הדוגמאות של Beam-CPP.

  1. משתמשים בפקודה git clone כדי לשכפל את מאגר GitHub:

    git clone https://github.com/GoogleCloudPlatform/dataflow-sample-applications.git
    
  2. עוברים לספריית האפליקציה:

    cd dataflow-sample-applications/beam-cpp-example
    

קוד הפייפליין

אתם יכולים להתאים אישית את קוד צינור הנתונים מהמדריך הזה. צינור הנתונים הזה משלים את המשימות הבאות:

  • יוצר באופן דינמי את כל המספרים השלמים בטווח קלט.
  • הפונקציה מריצה את המספרים השלמים דרך פונקציית C++ ומסננת ערכים לא תקינים.
  • כותבת את הערכים הלא תקינים לערוץ צדדי.
  • סופר את המקרים של כל זמן עצירה ומנרמל את התוצאות.
  • מדפיס את הפלט, מעצב את התוצאות וכותב אותן לקובץ טקסט.
  • יוצרת PCollection עם רכיב יחיד.
  • מעבד את הרכיב היחיד באמצעות פונקציית map ומעביר את התדירות PCollection כקלט צדדי.
  • מעבד את PCollection ומפיק פלט יחיד.

קובץ ההתחלה נראה כך:

#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements.  See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License.  You may obtain a copy of the License at
#
#    http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#


import argparse
import logging
import os
import sys


def run(argv):
  # Import here to avoid __main__ session pickling issues.
  import io
  import itertools
  import matplotlib.pyplot as plt
  import collatz

  import apache_beam as beam
  from apache_beam.io import restriction_trackers
  from apache_beam.options.pipeline_options import PipelineOptions

  class RangeSdf(beam.DoFn, beam.RestrictionProvider):
    """An SDF producing all the integers in the input range.

    This is preferable to beam.Create(range(...)) as it produces the integers
    dynamically rather than materializing them up front.  It is an SDF to do
    so with perfect dynamic sharding.
    """
    def initial_restriction(self, desired_range):
      start, stop = desired_range
      return restriction_trackers.OffsetRange(start, stop)

    def restriction_size(self, _, restriction):
      return restriction.size()

    def create_tracker(self, restriction):
      return restriction_trackers.OffsetRestrictionTracker(restriction)

    def process(self, _, active_range=beam.DoFn.RestrictionParam()):
      for i in itertools.count(active_range.current_restriction().start):
        if active_range.try_claim(i):
          yield i
        else:
          break

  class GenerateIntegers(beam.PTransform):
    def __init__(self, start, stop):
      self._start = start
      self._stop = stop

    def expand(self, p):
      return (
          p
          | beam.Create([(self._start, self._stop + 1)])
          | beam.ParDo(RangeSdf()))

  parser = argparse.ArgumentParser()
  parser.add_argument('--start', dest='start', type=int, default=1)
  parser.add_argument('--stop', dest='stop', type=int, default=10000)
  parser.add_argument('--output', default='./out.png')

  known_args, pipeline_args = parser.parse_known_args(argv)
  # Store this as a local to avoid capturing the full known_args.
  output_path = known_args.output

  with beam.Pipeline(options=PipelineOptions(pipeline_args)) as p:

    # Generate the integers from start to stop (inclusive).
    integers = p | GenerateIntegers(known_args.start, known_args.stop)

    # Run them through our C++ function, filtering bad records.
    # Requires apache beam 2.34 or later.
    stopping_times, bad_values = (
        integers
        | beam.Map(collatz.total_stopping_time).with_exception_handling(
            use_subprocess=True))

    # Write the bad values to a side channel.
    bad_values | 'WriteBadValues' >> beam.io.WriteToText(
        os.path.splitext(output_path)[0] + '-bad.txt')

    # Count the occurrence of each stopping time and normalize.
    total = known_args.stop - known_args.start + 1
    frequencies = (
        stopping_times
        | 'Aggregate' >> (beam.Map(lambda x: (x, 1)) | beam.CombinePerKey(sum))
        | 'Normalize' >> beam.MapTuple(lambda x, count: (x, count / total)))

    if known_args.stop <= 10:
      # Print out the results for debugging.
      frequencies | beam.Map(print)
    else:
      # Format and write them to a text file.
      (
          frequencies
          | 'Format' >> beam.MapTuple(lambda count, freq: f'{count}, {freq}')
          | beam.io.WriteToText(os.path.splitext(output_path)[0] + '.txt'))

    # Define some helper functions.
    def make_scatter_plot(xy):
      x, y = zip(*xy)
      plt.plot(x, y, '.')
      png_bytes = io.BytesIO()
      plt.savefig(png_bytes, format='png')
      png_bytes.seek(0)
      return png_bytes.read()

    def write_to_path(path, content):
      """Most Beam IOs write multiple elements to some kind of a container
      file (e.g. strings to lines of a text file, avro records to an avro file,
      etc.)  This function writes each element to its own file, given by path.
      """
      # Write to a temporary path and to a rename for fault tolerence.
      tmp_path = path + '.tmp'
      fs = beam.io.filesystems.FileSystems.get_filesystem(path)
      with fs.create(tmp_path) as fout:
        fout.write(content)
      fs.rename([tmp_path], [path])

    (
        p
        # Create a PCollection with a single element.
        | 'CreateSingleton' >> beam.Create([None])
        # Process the single element with a Map function, passing the frequency
        # PCollection as a side input.
        # This will cause the normally distributed frequency PCollection to be
        # colocated and processed as a single unit, producing a single output.
        | 'MakePlot' >> beam.Map(
            lambda _,
            data: make_scatter_plot(data),
            data=beam.pvalue.AsList(frequencies))
        # Pair this with the desired filename.
        |
        'PairWithFilename' >> beam.Map(lambda content: (output_path, content))
        # And actually write it out, using MapTuple to split the tuple into args.
        | 'WriteToOutput' >> beam.MapTuple(write_to_path))


if __name__ == '__main__':
  logging.getLogger().setLevel(logging.INFO)
  run(sys.argv)

הגדרת סביבת הפיתוח

  1. משתמשים ב-Apache Beam SDK ל-Python.

  2. מתקינים את ספריית GMP:

    apt-get install libgmp3-dev
    
  3. כדי להתקין את יחסי התלות, משתמשים בקובץ requirements.txt.

    pip install -r requirements.txt
    
  4. כדי ליצור את הקישורים של Python, מריצים את הפקודה הבאה.

    python setup.py build_ext --inplace
    

אתם יכולים להתאים אישית את קובץ requirements.txt מהמדריך הזה. קובץ ההתחלה כולל את יחסי התלות הבאים:

#
#    Licensed to the Apache Software Foundation (ASF) under one or more
#    contributor license agreements.  See the NOTICE file distributed with
#    this work for additional information regarding copyright ownership.
#    The ASF licenses this file to You under the Apache License, Version 2.0
#    (the "License"); you may not use this file except in compliance with
#    the License.  You may obtain a copy of the License at
#
#       http://www.apache.org/licenses/LICENSE-2.0
#
#    Unless required by applicable law or agreed to in writing, software
#    distributed under the License is distributed on an "AS IS" BASIS,
#    WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
#    See the License for the specific language governing permissions and
#    limitations under the License.
#

apache-beam[gcp]==2.46.0
cython==0.29.24
pyparsing==2.4.2
matplotlib==3.4.3

הפעלת הפייפליין באופן מקומי

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

כדי להריץ את צינור העיבוד באופן מקומי, משתמשים בפקודה הבאה. הפקודה הזו יוצרת תמונה בשם out.png.

python pipeline.py

יצירת Google Cloud המשאבים

בקטע הזה מוסבר איך ליצור את המשאבים הבאים:

  • קטגוריה של Cloud Storage לשימוש כמיקום אחסון זמני ומיקום פלט.
  • קונטיינר Docker לאריזת הקוד של צינור עיבוד הנתונים ויחסי התלות.

יצירת קטגוריה של Cloud Storage

מתחילים ביצירת קטגוריה של Cloud Storage באמצעות Google Cloud CLI. הבאקט הזה משמש כמיקום אחסון זמני על ידי צינור Dataflow.

כדי ליצור את הקטגוריה, משתמשים בפקודה gcloud storage buckets create:

gcloud storage buckets create gs://BUCKET_NAME --location=LOCATION

מחליפים את מה שכתוב בשדות הבאים:

יצירה ובנייה של תמונת קונטיינר

אתם יכולים להתאים אישית את קובץ ה-Dockerfile מהמדריך הזה. קובץ ההתחלה נראה כך:

FROM apache/beam_python3.9_sdk:2.46.0

# Install a C++ library.
RUN apt-get update
RUN apt-get install -y libgmp3-dev

# Install Python dependencies.
COPY requirements.txt requirements.txt
RUN pip install -r requirements.txt

# Install the code and some python bindings.
COPY pipeline.py pipeline.py
COPY collatz.pyx collatz.pyx
COPY setup.py setup.py
RUN python setup.py install

קובץ Docker הזה מכיל את הפקודות FROM,‏ COPY ו-RUN, שאפשר לקרוא עליהן בהפניה לקובץ Docker.

  1. כדי להעלות פריטי מידע שנוצרו בתהליך פיתוח (Artifact), צריך ליצור מאגר Artifact Registry. כל מאגר יכול להכיל פריטי מידע שנוצרו בתהליך פיתוח (Artifact) בפורמט נתמך אחד בלבד.

    כל התוכן במאגר מוצפן באמצעות Google-owned and Google-managed encryption keys או מפתחות הצפנה שמנוהלים על ידי הלקוח. ב-Artifact Registry נעשה שימוש ב-Google-owned and Google-managed encryption keys כברירת מחדל, ולא נדרשת הגדרה לאפשרות הזו.

    צריכה להיות לכם לפחות הרשאת כתיבה ב-Artifact Registry למאגר.

    מריצים את הפקודה הבאה כדי ליצור מאגר חדש. הפקודה משתמשת בדגל --async וחוזרת מיידית, בלי להמתין לסיום הפעולה.

    gcloud artifacts repositories create REPOSITORY \
       --repository-format=docker \
       --location=LOCATION \
       --async
    

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

  2. יוצרים את קובץ Dockerfile.

    כדי שחבילות יהיו חלק מהמאגר של Apache Beam, צריך לציין אותן כחלק מקובץ requirements.txt. חשוב לוודא שלא מציינים את apache-beam כחלק מקובץ requirements.txt. הקונטיינר של Apache Beam כבר כולל את apache-beam.

  3. כדי לדחוף או לשלוף תמונות, צריך להגדיר את Docker לאימות בקשות ל-Artifact Registry. כדי להגדיר אימות למאגרי Docker, מריצים את הפקודה הבאה:

    gcloud auth configure-docker LOCATION-docker.pkg.dev
    

    הפקודה מעדכנת את ההגדרה של Docker. מעכשיו אפשר להתחבר אל Artifact Registry בפרויקט Google Cloud כדי להעלות תמונות.

  4. יוצרים את קובץ האימג' של Docker באמצעות Dockerfile עם Cloud Build.

    מעדכנים את הנתיב בפקודה הבאה כך שיתאים לקובץ ה-Dockerfile שיצרתם. הפקודה הזו יוצרת את הקובץ ומעבירה אותו בדחיפה למאגר שלכם ב-Artifact Registry.

    gcloud builds submit --tag LOCATION-docker.pkg.dev/PROJECT_ID/REPOSITORY/dataflow/cpp_beam_container:latest .
    

אריזת הקוד והתלויות בקונטיינר של Docker

  1. כדי להריץ את צינור העיבוד הזה בסביבה מבוזרת, צריך לארוז את הקוד ואת התלות בקונטיינר של Docker.

    docker build . -t cpp_beam_container
    
  2. אחרי שאורזים את הקוד ואת התלות, אפשר להריץ את צינור הנתונים באופן מקומי כדי לבדוק אותו.

    python pipeline.py \
       --runner=PortableRunner \
       --job_endpoint=embed \
       --environment_type=DOCKER \
       --environment_config="docker.io/library/cpp_beam_container"
    

    הפקודה הזו כותבת את הפלט בתוך תמונת Docker. כדי לראות את הפלט, מריצים את צינור העיבוד עם --output וכותבים את הפלט לקטגוריית אחסון ב-Cloud Storage. לדוגמה, מריצים את הפקודה הבאה.

    python pipeline.py \
       --runner=PortableRunner \
       --job_endpoint=embed \
       --environment_type=DOCKER \
       --environment_config="docker.io/library/cpp_beam_container" \
       --output=gs://BUCKET_NAME/out.png
    

הרצת צינור עיבוד הנתונים

עכשיו אפשר להריץ את צינור עיבוד הנתונים של Apache Beam ב-Dataflow. לשם כך, צריך להפנות לקובץ עם קוד צינור עיבוד הנתונים ולהעביר את הפרמטרים שנדרשים לצינור עיבוד הנתונים.

במעטפת או בטרמינל, מריצים את צינור עיבוד הנתונים באמצעות Dataflow Runner.

python pipeline.py \
    --runner=DataflowRunner \
    --project=PROJECT_ID \
    --region=REGION \
    --temp_location=gs://BUCKET_NAME/tmp \
    --sdk_container_image="LOCATION-docker.pkg.dev/PROJECT_ID/REPOSITORY/dataflow/cpp_beam_container:latest" \
    --experiment=use_runner_v2 \
    --output=gs://BUCKET_NAME/out.png

אחרי שמריצים את הפקודה להפעלת צינור הנתונים, Dataflow מחזיר מזהה משימה עם סטטוס המשימה Queued (בהמתנה). יכול להיות שיעברו כמה דקות עד שהסטטוס של העבודה יהיה Running ותהיה לכם גישה לתרשים העבודה.

צפייה בתוצאות

הצגת נתונים שנכתבו לקטגוריה ב-Cloud Storage. משתמשים בפקודה gcloud storage ls כדי להציג את התוכן ברמה העליונה של הקטגוריה:

gcloud storage ls gs://BUCKET_NAME

אם הפקודה מסתיימת בלי שגיאות, היא מחזירה הודעה שדומה לזו:

gs://BUCKET_NAME/out.png

הסרת המשאבים

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

מחיקת הפרויקט

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

  1. במסוף Google Cloud , נכנסים לדף Manage resources.

    כניסה לדף Manage resources

  2. ברשימת הפרויקטים, בוחרים את הפרויקט שרוצים למחוק ולוחצים על Delete.
  3. כדי למחוק את הפרויקט, כותבים את מזהה הפרויקט בתיבת הדו-שיח ולוחצים על Shut down.

מחיקת המשאבים הבודדים

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

מחיקת משאבי הפרויקט Google Cloud

  1. מחיקת מאגר Artifact Registry.

    gcloud artifacts repositories delete REPOSITORY \
       --location=LOCATION --async
    
  2. מוחקים את הקטגוריה של Cloud Storage ואת האובייקטים שבה. השימוש בדלי הזה לבדו לא כרוך בחיובים.

    gcloud storage rm gs://BUCKET_NAME --recursive
    

ביטול פרטי כניסה

  1. מבטלים את התפקידים שהוקצו לחשבון השירות של העובד (worker) שמנוהל על ידי משתמש. מריצים את הפקודה הבאה לכל אחד מהתפקידים הבאים ב-IAM:

    • roles/dataflow.admin
    • roles/dataflow.worker
    • roles/storage.objectAdmin
    • roles/artifactregistry.reader
    gcloud projects remove-iam-policy-binding PROJECT_ID \
      --member=serviceAccount:parallelpipeline@PROJECT_ID.iam.gserviceaccount.com \
      --role=SERVICE_ACCOUNT_ROLE
  2. אם תרצו, תוכלו לבטל את פרטי הכניסה שיצרתם ולמחוק את הקובץ המקומי של פרטי הכניסה.

    gcloud auth application-default revoke
  3. אם רוצים, מבטלים את פרטי הכניסה של ה-CLI של gcloud.

    gcloud auth revoke

המאמרים הבאים