הרצת משימה של ספירת מילים ב-Hadoop באשכול של Managed Service for Apache Spark

Managed Airflow (דור 3) | Managed Airflow (דור 2) | Managed Airflow (דור 1 מדור קודם)

במדריך הזה נסביר איך להשתמש ב-Managed Service for Apache Airflow כדי ליצור גרף אציקלי מכוון (DAG) של Apache Airflow שמריץ משימת ספירת מילים של Apache Hadoop באשכול Managed Service for Apache Spark.

מטרות

  1. ניגשים לסביבת Managed Airflow ומשתמשים בממשק המשתמש של Airflow.
  2. יצירה של משתני סביבה ב-Airflow וצפייה בהם.
  3. יצירה והרצה של DAG שכולל את המשימות הבאות:
    1. יוצר אשכול של Managed Service for Apache Spark.
    2. מריץ משימה של ספירת מילים ב-Apache Hadoop באשכול.
    3. התוצאות של ספירת המילים מועברות לקטגוריה של Cloud Storage.
    4. מחיקת האשכול.

עלויות

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

  • Managed Airflow
  • Managed Service for Apache Spark
  • Cloud Storage

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

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

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

  • מוודאים שממשקי ה-API הבאים מופעלים בפרויקט:

    המסוף

    מפעילים את Managed Service for Apache Spark ואת Cloud Storage APIs.

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

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

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

    gcloud

    מפעילים את Managed Service for Apache Spark, Cloud Storage APIs:

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

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

    gcloud services enable dataproc.googleapis.com storage-component.googleapis.com

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

  • שימו לב לנתיב של הקטגוריה שיצרתם, לדוגמה: gs://example-bucket. בהמשך המדריך הזה, תגדירו משתנה Airflow לנתיב הזה ותשתמשו במשתנה בדוגמה ל-DAG.

  • יצירת סביבת Managed Airflow עם פרמטרים שמוגדרים כברירת מחדל. מחכים עד ליצירת הסביבה. בסיום, סימן וי ירוק מופיע מימין לשם הסביבה.

  • שימו לב לאזור שבו יצרתם את הסביבה, למשל us-central. תגדירו משתנה Airflow לאזור הזה ותשתמשו בו ב-DAG לדוגמה כדי להפעיל אשכול Managed Service for Apache Spark באותו אזור.

הגדרת משתני Airflow

מגדירים את המשתנים של Airflow כדי להשתמש בהם בהמשך בדוגמה של DAG. לדוגמה, אפשר להגדיר משתני Airflow בממשק המשתמש של Airflow.

משתנה Airflow ערך
gcp_project מזהה הפרויקט שבו אתם משתמשים במדריך הזה, למשל example-project.
gcs_bucket ה-URI של קטגוריה של Cloud Storage שיצרתם לצורך המדריך הזה, כמו gs://example-bucket.
gce_region האזור שבו יצרתם את הסביבה, למשל us-central1. זהו האזור שבו ייווצר האשכול של Managed Service for Apache Spark.

הצגת זרימת העבודה לדוגמה

‫DAG ב-Airflow הוא אוסף של משימות מאורגנות שרוצים לתזמן ולהריץ. גרפים מכווני מחזור (DAG) מוגדרים בקובצי Python רגילים. הקוד שמוצג ב-hadoop_tutorial.py הוא קוד תהליך העבודה.

"""Example Airflow DAG that creates a Cloud Dataproc cluster, runs the Hadoop
wordcount example, and deletes the cluster.

This DAG relies on three Airflow variables
https://airflow.apache.org/docs/apache-airflow/stable/concepts/variables.html
* gcp_project - Google Cloud Project to use for the Cloud Dataproc cluster.
* gce_region - Google Compute Engine region where Cloud Dataproc cluster should be
  created.
* gcs_bucket - Google Cloud Storage bucket to use for result of Hadoop job.
  See https://cloud.google.com/storage/docs/creating-buckets for creating a
  bucket.
"""

import datetime
import os

from airflow import models
from airflow.providers.google.cloud.operators import dataproc
from airflow.utils import trigger_rule

# Output file for Cloud Dataproc job.
# If you are running Airflow in more than one time zone
# see https://airflow.apache.org/docs/apache-airflow/stable/timezone.html
# for best practices
output_file = (
    os.path.join(
        "{{ var.value.gcs_bucket }}",
        "wordcount",
        datetime.datetime.now().strftime("%Y%m%d-%H%M%S"),
    )
    + os.sep
)
# Path to Hadoop wordcount example available on every Dataproc cluster.
WORDCOUNT_JAR = "file:///usr/lib/hadoop-mapreduce/hadoop-mapreduce-examples.jar"
# Arguments to pass to Cloud Dataproc job.
input_file = "gs://pub/shakespeare/rose.txt"
wordcount_args = ["wordcount", input_file, output_file]

HADOOP_JOB = {
    "reference": {"project_id": "{{ var.value.gcp_project }}"},
    "placement": {"cluster_name": "composer-hadoop-tutorial-cluster-{{ ds_nodash }}"},
    "hadoop_job": {
        "main_jar_file_uri": WORDCOUNT_JAR,
        "args": wordcount_args,
    },
}

CLUSTER_CONFIG = {
    "master_config": {"num_instances": 1, "machine_type_uri": "n1-standard-2"},
    "worker_config": {"num_instances": 2, "machine_type_uri": "n1-standard-2"},
}

yesterday = datetime.datetime.combine(
    datetime.datetime.today() - datetime.timedelta(1), datetime.datetime.min.time()
)

default_dag_args = {
    # Setting start date as yesterday starts the DAG immediately when it is
    # detected in the Cloud Storage bucket.
    "start_date": yesterday,
    # To email on failure or retry set 'email' arg to your email and enable
    # emailing here.
    "email_on_failure": False,
    "email_on_retry": False,
    # If a task fails, retry it once after waiting at least 5 minutes
    "retries": 1,
    "retry_delay": datetime.timedelta(minutes=5),
    "project_id": "{{ var.value.gcp_project }}",
    "region": "{{ var.value.gce_region }}",
}


with models.DAG(
    "composer_hadoop_tutorial",
    # Continue to run DAG once per day
    schedule_interval=datetime.timedelta(days=1),
    default_args=default_dag_args,
) as dag:

    # Create a Cloud Dataproc cluster.
    create_dataproc_cluster = dataproc.DataprocCreateClusterOperator(
        task_id="create_dataproc_cluster",
        # Give the cluster a unique name by appending the date scheduled.
        # See https://airflow.apache.org/docs/apache-airflow/stable/macros-ref.html
        cluster_name="composer-hadoop-tutorial-cluster-{{ ds_nodash }}",
        cluster_config=CLUSTER_CONFIG,
        region="{{ var.value.gce_region }}",
    )

    # Run the Hadoop wordcount example installed on the Cloud Dataproc cluster
    # master node.
    run_dataproc_hadoop = dataproc.DataprocSubmitJobOperator(
        task_id="run_dataproc_hadoop", job=HADOOP_JOB
    )

    # Delete Cloud Dataproc cluster.
    delete_dataproc_cluster = dataproc.DataprocDeleteClusterOperator(
        task_id="delete_dataproc_cluster",
        cluster_name="composer-hadoop-tutorial-cluster-{{ ds_nodash }}",
        region="{{ var.value.gce_region }}",
        # Setting trigger_rule to ALL_DONE causes the cluster to be deleted
        # even if the Dataproc job fails.
        trigger_rule=trigger_rule.TriggerRule.ALL_DONE,
    )

    # Define DAG dependencies.
    create_dataproc_cluster >> run_dataproc_hadoop >> delete_dataproc_cluster

אופרטורים

כדי לתזמן את שלוש המשימות בתהליך העבודה לדוגמה, ה-DAG מייבא את שלושת האופרטורים הבאים של Airflow:

  • DataprocClusterCreateOperator: יצירת אשכול של Managed Service for Apache Spark.

  • DataProcHadoopOperator: שולח משימת ספירת מילים של Hadoop וכותב את התוצאות לקטגוריה ב-Cloud Storage.

  • DataprocClusterDeleteOperator: מחיקת האשכול כדי להימנע מחיובים שוטפים ב-Compute Engine.

תלויות

אתם מארגנים את המשימות שאתם רוצים להריץ באופן שמשקף את הקשרים ויחסי התלות ביניהן. המשימות ב-DAG הזה פועלות ברצף.

# Define DAG dependencies.
create_dataproc_cluster >> run_dataproc_hadoop >> delete_dataproc_cluster

קביעת מועד

השם של ה-DAG הוא composer_hadoop_tutorial, והוא מופעל פעם ביום. כי הערך של start_date שמועבר אל default_dag_args מוגדר כ-yesterday. מערכת Managed Airflow מתזמנת את תהליך העבודה כך שיתחיל מיד אחרי שה-DAG מועלה אל דלי האחסון של הסביבה.

with models.DAG(
    "composer_hadoop_tutorial",
    # Continue to run DAG once per day
    schedule_interval=datetime.timedelta(days=1),
    default_args=default_dag_args,
) as dag:

העלאת ה-DAG לקטגוריה של הסביבה

‫Managed Airflow מאחסן DAGs בתיקייה /dags בדלי של הסביבה.

כדי להעלות את ה-DAG:

  1. במחשב המקומי, שומרים את הקובץ hadoop_tutorial.py.

  2. נכנסים לדף Environments במסוף Google Cloud .

    מעבר אל Environments

  3. ברשימת הסביבות, בעמודה DAGs folder של הסביבה, לוחצים על הקישור DAGs.

  4. לוחצים על העלאת קבצים.

  5. בוחרים באפשרות hadoop_tutorial.py במחשב המקומי ולוחצים על פתיחה.

‫Managed Airflow מוסיף את ה-DAG ל-Airflow ומתזמן את ה-DAG באופן אוטומטי. השינויים ב-DAG מתרחשים תוך 3-5 דקות.

עיון בהפעלות של DAG

צפייה בסטטוס המשימה

כשמעלים את קובץ ה-DAG לתיקייה dags/ ב-Cloud Storage,‏ Managed Airflow מנתח את הקובץ. אם הפעולה תושלם בהצלחה, השם של תהליך העבודה יופיע ברשימת ה-DAG, ותהליך העבודה יתווסף לתור להפעלה מיידית.

  1. כדי לראות את סטטוס המשימה, עוברים לממשק האינטרנט של Airflow ולוחצים על DAGs בסרגל הכלים.

  2. כדי לפתוח את דף הפרטים של DAG, לוחצים על composer_hadoop_tutorial. בדף הזה מוצג ייצוג גרפי של משימות ויחסי תלות בתהליך העבודה.

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

הוספה חוזרת של תהליך העבודה לתור

כדי להריץ שוב את תהליך העבודה מתצוגת התרשים:

  1. בתצוגת הגרף בממשק המשתמש של Airflow, לוחצים על הגרפיקה create_dataproc_cluster.
  2. כדי לאפס את שלוש המשימות, לוחצים על ניקוי ואז על אישור.
  3. לוחצים שוב על create_dataproc_cluster בתצוגת הגרף.
  4. כדי להוסיף שוב את תהליך העבודה לתור, לוחצים על הפעלה.

צפייה בתוצאות של משימות

אפשר גם לבדוק את הסטטוס והתוצאות של זרימת העבודה בדפים הבאים במסוף:composer_hadoop_tutorial Google Cloud

  • Managed Service for Apache Spark Clusters: כדי לעקוב אחרי יצירה ומחיקה של אשכולות. שימו לב שהאשכול שנוצר על ידי תהליך העבודה הוא זמני: הוא קיים רק למשך תהליך העבודה ונמחק כחלק מהמשימה האחרונה בתהליך העבודה.

    מעבר אל Managed Service for Apache Spark Clusters

  • Managed Service for Apache Spark Jobs: כדי להציג או לנטר את משימת ספירת המילים של Apache Hadoop. לוחצים על מזהה משימה כדי לראות את הפלט של יומן העבודה.

    מעבר אל Managed Service for Apache Spark Jobs

  • Cloud Storage Browser: כדי לראות את התוצאות של ספירת המילים בתיקייה wordcount בקטגוריה של Cloud Storage שיצרתם לצורך המדריך הזה.

    כניסה לדף Cloud Storage Browser

הסרת המשאבים

מוחקים את המשאבים שבהם השתמשתם במדריך הזה:

  1. מחיקת סביבת Managed Airflow, כולל מחיקה ידנית של קטגוריית הסביבה.

  2. מחיקת קטגוריית Cloud Storage שבה מאוחסנות התוצאות של משימת ספירת המילים ב-Hadoop.