שימוש ב-KubernetesPodOperator

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

בדף הזה מוסבר איך להשתמש ב-KubernetesPodOperator כדי לפרוס Kubernetes Pods מ-Managed Service for Apache Airflow אל אשכול Google Kubernetes Engine שמהווה חלק מהסביבה של Managed Service for Apache Airflow.

האופרטור KubernetesPodOperator מפעיל קבוצות Pod של Kubernetes באשכול שבסביבה שלכם. לעומת זאת, אופרטורים של Google Kubernetes Engine מריצים Pods של Kubernetes באשכול שצוין, שיכול להיות אשכול נפרד שלא קשור לסביבה שלכם. אפשר גם ליצור ולמחוק אשכולות באמצעות אופרטורים של Google Kubernetes Engine.

‫KubernetesPodOperator היא אפשרות טובה אם אתם צריכים:

  • יחסי תלות מותאמים אישית של Python שלא זמינים דרך מאגר PyPI הציבורי.
  • תלויות בינאריות שלא זמינות בתמונת ה-worker של Managed Airflow.

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

כדאי לעיין ברשימת ההבדלים הבאה בין KubernetesPodOperator ב-Managed Airflow (דור 3) לבין Managed Airflow (דור 2) ולוודא שקובצי ה-DAG שלכם תואמים:

  • אי אפשר ליצור מרחבי שמות בהתאמה אישית ב-Managed Airflow (דור 3). פודים תמיד פועלים במרחב השמות composer-user-workloads, גם אם מצוין מרחב שמות אחר. ל-Pods במרחב השמות הזה יש גישה למשאבים של הפרויקט ולרשת ה-VPC (אם היא מופעלת) ללא צורך בהגדרה נוספת.

  • אי אפשר להפעיל כמה קונטיינרים נוספים של sidecar ב-Managed Airflow (דור 3). אפשר להריץ קונטיינר sidecar יחיד אם השם שלו הוא airflow-xcom-sidecar.

  • אי אפשר ליצור סודות ומפות ConfigMap של Kubernetes באמצעות Kubernetes API. במקום זאת, Managed Airflow מספק פקודות של Google Cloud CLI, משאבי Terraform ו-Cloud Composer API לניהול סודות ו-ConfigMaps של Kubernetes. למידע נוסף, ראו שימוש ב-Kubernetes Secrets וב-ConfigMaps.

  • אי אפשר לפרוס עומסי עבודה מותאמים אישית ב-Managed Airflow (דור 3). אפשר לשנות רק את Kubernetes Secrets ואת ConfigMaps, אבל אי אפשר לבצע שינויים אחרים בהגדרות.

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

  • בדומה ל-Managed Airflow (דור 2), אי אפשר להגדיר שיוך של Pod. אם רוצים להשתמש בהעדפה של Pod, אפשר להשתמש באופרטורים של GKE כדי להפעיל Pods באשכול אחר.

מידע על KubernetesPodOperator ב-Managed Airflow (דור 3)

בקטע הזה מוסבר איך KubernetesPodOperator פועל ב-Managed Airflow (דור 3).

שימוש במשאבים

ב-Managed Airflow (דור 3), האשכול של הסביבה מתרחב אוטומטית. עומסי עבודה נוספים שאתם מריצים באמצעות KubernetesPodOperator מתרחבים באופן עצמאי מהסביבה שלכם. הסביבה שלכם לא מושפעת מהעלייה בביקוש למשאבים, אבל גודל האשכול בסביבה משתנה בהתאם לביקוש למשאבים.

התמחור של עומסי העבודה הנוספים שמופעלים באשכול של הסביבה מבוסס על מודל התמחור של Managed Airflow (דור 3) ומשתמש במק"טים של Managed Airflow (דור 3).

‫Managed Airflow (דור 3) משתמש באשכולות Autopilot, שכוללים את המושג compute classes:

  • ב-Managed Airflow יש תמיכה רק בסוג המחשוב general-purpose.

  • כברירת מחדל, אם לא נבחרה כיתה, המערכת מניחה שמדובר בכיתה general-purpose כשיוצרים תרמילים באמצעות KubernetesPodOperator.

  • כל מחלקה משויכת למאפיינים ספציפיים ולמגבלות משאבים. אפשר לקרוא עליהם במסמכי Autopilot. לדוגמה, ל-Pods שפועלים בתוך המחלקה general-purpose יכולים להיות עד ‎110 GiB של זיכרון.

גישה למשאבים של הפרויקט

ב-Managed Airflow (דור 3), האשכול של הסביבה נמצא בפרויקט הדייר, ואי אפשר להגדיר אותו. הפעלת ה-Pods מתבצעת באשכול של הסביבה, במרחב שמות מבודד.

ב-Managed Airflow (דור 3), הפודים תמיד פועלים במרחב השמות composer-user-workloads, גם אם צוין מרחב שמות אחר. לפודים במרחב השמות הזה יש גישה למשאבים בפרויקט ולרשת ה-VPC (אם היא מופעלת) ללא צורך בהגדרה נוספת. Google Cloud חשבון השירות של הסביבה משמש לגישה למשאבים האלה. אי אפשר לציין חשבון שירות אחר.

הגדרה מינימלית

כדי ליצור KubernetesPodOperator, נדרשים רק הפרמטרים name,‏ image ו-task_id של Pod. ‫/home/airflow/composer_kube_config מכיל פרטי כניסה לאימות ב-GKE.

kubernetes_min_pod = KubernetesPodOperator(
    # The ID specified for the task.
    task_id="pod-ex-minimum",
    # Name of task you want to run, used to generate Pod ID.
    name="pod-ex-minimum",
    # Entrypoint of the container, if not specified the Docker container's
    # entrypoint is used. The cmds parameter is templated.
    cmds=["echo"],
    # The namespace to run within Kubernetes. In Composer 2 environments
    # after December 2022, the default namespace is
    # `composer-user-workloads`. Always use the
    # `composer-user-workloads` namespace with Composer 3.
    namespace="composer-user-workloads",
    # Docker image specified. Defaults to hub.docker.com, but any fully
    # qualified URLs will point to a custom repository. Supports private
    # gcr.io images if the Composer Environment is under the same
    # project-id as the gcr.io images and the service account that Composer
    # uses has permission to access the Google Container Registry
    # (the default service account has permission)
    image="marketplace.gcr.io/google/ubuntu2204",
    # Specifies path to kubernetes config. The config_file is templated.
    config_file="/home/airflow/composer_kube_config",
    # Identifier of connection that should be used
    kubernetes_conn_id="kubernetes_default",
)

שלבי הגדרת תצורה נוספים

בדוגמה הזו מוצגים פרמטרים נוספים שאפשר להגדיר ב-KubernetesPodOperator.

מידע נוסף זמין במקורות המידע הבאים:

kubernetes_full_pod = KubernetesPodOperator(
    task_id="ex-all-configs",
    name="pi",
    namespace="composer-user-workloads",
    image="perl:5.34.0",
    # Entrypoint of the container, if not specified the Docker container's
    # entrypoint is used. The cmds parameter is templated.
    cmds=["perl"],
    # Arguments to the entrypoint. The Docker image's CMD is used if this
    # is not provided. The arguments parameter is templated.
    arguments=["-Mbignum=bpi", "-wle", "print bpi(2000)"],
    # The secrets to pass to Pod, the Pod will fail to create if the
    # secrets you specify in a Secret object do not exist in Kubernetes.
    secrets=[],
    # Labels to apply to the Pod.
    labels={"pod-label": "label-name"},
    # Timeout to start up the Pod, default is 600.
    startup_timeout_seconds=600,
    # The environment variables to be initialized in the container.
    # The env_vars parameter is templated.
    env_vars={"EXAMPLE_VAR": "/example/value"},
    # If true, logs stdout output of container. Defaults to True.
    get_logs=True,
    # Determines when to pull a fresh image, if 'IfNotPresent' will cause
    # the Kubelet to skip pulling an image if it already exists. If you
    # want to always pull a new image, set it to 'Always'.
    image_pull_policy="Always",
    # Annotations are non-identifying metadata you can attach to the Pod.
    # Can be a large range of data, and can include characters that are not
    # permitted by labels.
    annotations={"key1": "value1"},
    # Optional resource specifications for Pod, this will allow you to
    # set both cpu and memory limits and requirements.
    # Prior to Airflow 2.3 and the cncf providers package 5.0.0
    # resources were passed as a dictionary. This change was made in
    # https://github.com/apache/airflow/pull/27197
    # Additionally, "memory" and "cpu" were previously named
    # "limit_memory" and "limit_cpu"
    # resources={'limit_memory': "250M", 'limit_cpu': "100m"},
    container_resources=k8s_models.V1ResourceRequirements(
        requests={"cpu": "1000m", "memory": "10G", "ephemeral-storage": "10G"},
        limits={"cpu": "1000m", "memory": "10G", "ephemeral-storage": "10G"},
    ),
    # Specifies path to kubernetes config. The config_file is templated.
    config_file="/home/airflow/composer_kube_config",
    # If true, the content of /airflow/xcom/return.json from container will
    # also be pushed to an XCom when the container ends.
    do_xcom_push=False,
    # List of Volume objects to pass to the Pod.
    volumes=[],
    # List of VolumeMount objects to pass to the Pod.
    volume_mounts=[],
    # Identifier of connection that should be used
    kubernetes_conn_id="kubernetes_default",
    # Affinity determines which nodes the Pod can run on based on the
    # config. For more information see:
    # https://kubernetes.io/docs/concepts/configuration/assign-pod-node/
    # Pod affinity with the KubernetesPodOperator
    # is not supported with Composer 2
    # instead, create a cluster and use the GKEStartPodOperator
    # https://cloud.google.com/composer/docs/using-gke-operator
    affinity={},
)

שימוש בתבניות Jinja

‫Airflow תומך בתבניות Jinja ב-DAG.

צריך להצהיר על פרמטרים נדרשים של Airflow ‏ (task_id,‏ name ו-image) באמצעות האופרטור. כמו שרואים בדוגמה הבאה, אפשר להשתמש ב-Jinja כדי ליצור תבנית לכל הפרמטרים האחרים, כולל cmds,‏ arguments,‏ env_vars ו-config_file.

הפרמטר env_vars בדוגמה מוגדר ממשתנה Airflow בשם my_value. ערך ה-DAG בדוגמה מתקבל ממשתנה התבנית vars ב-Airflow. ל-Airflow יש עוד משתנים שמאפשרים גישה לסוגים שונים של מידע. לדוגמה, אפשר להשתמש במשתנה התבנית conf כדי לגשת לערכים של אפשרויות ההגדרה של Airflow. מידע נוסף ורשימת המשתנים שזמינים ב-Airflow מופיעים בחומר העזר בנושא תבניות במסמכי התיעוד של Airflow.

בלי לשנות את ה-DAG או ליצור את המשתנה env_vars, המשימה ex-kube-templates בדוגמה נכשלת כי המשתנה לא קיים. יוצרים את המשתנה הזה בממשק המשתמש של Airflow או באמצעות Google Cloud CLI:

ממשק המשתמש של Airflow

  1. עוברים אל ממשק המשתמש של Airflow.

  2. בסרגל הכלים, בוחרים באפשרות אדמין > משתנים.

  3. בדף List Variable (משתנה רשימה), לוחצים על Add a new record (הוספת רשומה חדשה).

  4. בדף Add Variable (הוספת משתנה), מזינים את הפרטים הבאים:

    • מקש:my_value
    • ערך: example_value
  5. לוחצים על Save.

gcloud

מזינים את הפקודה הבאה:

gcloud composer environments run ENVIRONMENT \
    --location LOCATION \
    variables set -- \
    my_value example_value

מחליפים את:

  • ENVIRONMENT בשם הסביבה.
  • LOCATION עם האזור שבו הסביבה ממוקמת.

בדוגמה הבאה אפשר לראות איך משתמשים בתבניות Jinja עם KubernetesPodOperator:

kubernetes_template_ex = KubernetesPodOperator(
    task_id="ex-kube-templates",
    name="ex-kube-templates",
    namespace="composer-user-workloads",
    image="bash",
    # All parameters below can be templated with Jinja. For more information
    # and the list of variables available in Airflow, see
    # the Airflow templates reference:
    # https://airflow.apache.org/docs/apache-airflow/stable/templates-ref.html
    # Entrypoint of the container, if not specified the Docker container's
    # entrypoint is used. The cmds parameter is templated.
    cmds=["echo"],
    # DS in Jinja is the execution date as YYYY-MM-DD, this Docker image
    # will echo the execution date. Arguments to the entrypoint. The Docker
    # image's CMD is used if this is not provided. The arguments parameter
    # is templated.
    arguments=["{{ ds }}"],
    # The var template variable allows you to access variables defined in
    # Airflow UI. In this case we are getting the value of my_value and
    # setting the environment variable `MY_VALUE`. The pod will fail if
    # `my_value` is not set in the Airflow UI. The env_vars parameter
    # is templated.
    env_vars={"MY_VALUE": "{{ var.value.my_value }}"},
    # Specifies path to Kubernetes config. The config_file is templated.
    config_file="/home/airflow/composer_kube_config",
    # Identifier of connection that should be used
    kubernetes_conn_id="kubernetes_default",
)

שימוש בסודות וב-ConfigMaps של Kubernetes

סוד ב-Kubernetes הוא אובייקט שמכיל מידע רגיש. ‫ConfigMap ב-Kubernetes הוא אובייקט שמכיל נתונים לא סודיים בצמדי מפתח-ערך.

ב-Managed Airflow (דור 3), אפשר ליצור סודות ו-ConfigMaps באמצעות Google Cloud CLI,‏ API או Terraform, ואז לגשת אליהם מ-KubernetesPodOperator:

  • ב-Google Cloud CLI וב-API, אתם מספקים קובץ הגדרות בפורמט YAML.
  • ב-Terraform, מגדירים סודות ו-ConfigMaps כמשאבים נפרדים בקובצי התצורה של Terraform.

מידע על קובצי הגדרות ב-YAML

כשיוצרים Kubernetes Secret או ConfigMap באמצעות Google Cloud CLI ו-API, צריך לספק קובץ בפורמט YAML. הפורמט של הקובץ הזה צריך להיות זהה לפורמט שבו משתמשים ב-Kubernetes Secrets וב-ConfigMaps. בתיעוד של Kubernetes יש הרבה דוגמאות קוד של ConfigMap ו-Secrets. כדי להתחיל, אפשר לעיין בדף הפצת פרטי כניסה בצורה מאובטחת באמצעות סודות ובמאמר בנושא ConfigMaps.

כמו ב-Kubernetes Secrets, צריך להשתמש בייצוג base64 כשמגדירים ערכים ב-Secrets.

כדי לקודד ערך, אפשר להשתמש בפקודה הבאה (זו אחת מתוך הרבה דרכים לקבל ערך מקודד בפורמט Base64):

echo "postgresql+psycopg2://root:example-password@127.0.0.1:3306/example-db" -n | base64

פלט:

cG9zdGdyZXNxbCtwc3ljb3BnMjovL3Jvb3Q6ZXhhbXBsZS1wYXNzd29yZEAxMjcuMC4wLjE6MzMwNi9leGFtcGxlLWRiIC1uCg==

בהמשך המדריך הזה נשתמש בשתי הדוגמאות הבאות של קובצי YAML. קובץ הגדרות לדוגמה ב-YAML בשביל סוד של Kubernetes:

apiVersion: v1
kind: Secret
metadata:
  name: airflow-secrets
data:
  sql_alchemy_conn: cG9zdGdyZXNxbCtwc3ljb3BnMjovL3Jvb3Q6ZXhhbXBsZS1wYXNzd29yZEAxMjcuMC4wLjE6MzMwNi9leGFtcGxlLWRiIC1uCg==

דוגמה נוספת שמראה איך לכלול קבצים. כמו בדוגמה הקודמת, קודם מקודדים את התוכן של קובץ (cat ./key.json | base64), ואז מספקים את הערך הזה בקובץ ה-YAML:

apiVersion: v1
kind: Secret
metadata:
  name: service-account
data:
  service-account.json: |
    ewogICJ0eXBl...mdzZXJ2aWNlYWNjb3VudC5jb20iCn0K

קובץ הגדרות לדוגמה ב-YAML עבור ConfigMap. אין צורך להשתמש בייצוג base64 ב-ConfigMaps:

apiVersion: v1
kind: ConfigMap
metadata:
  name: example-configmap
data:
  example_key: example_value

ניהול סודות ב-Kubernetes

gcloud

יצירת סוד

כדי ליצור סוד של Kubernetes, מריצים את הפקודה הבאה:

gcloud beta composer environments user-workloads-secrets create \
  --environment ENVIRONMENT_NAME \
  --location LOCATION \
  --secret-file-path SECRET_FILE

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

  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.
  • SECRET_FILE: הנתיב לקובץ YAML מקומי שמכיל את ההגדרה של הסוד.

דוגמה:

gcloud beta composer environments user-workloads-secrets create \
  --environment example-environment \
  --location us-central1 \
  --secret-file-path ./secrets/example-secret.yaml

עדכון סוד

כדי לעדכן סוד של Kubernetes, מריצים את הפקודה הבאה. השם של הסוד יילקח מקובץ ה-YAML שצוין, והתוכן של הסוד יוחלף.

gcloud beta composer environments user-workloads-secrets update \
  --environment ENVIRONMENT_NAME \
  --location LOCATION \
  --secret-file-path SECRET_FILE

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

  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.
  • SECRET_FILE: הנתיב לקובץ YAML מקומי שמכיל את ההגדרה של הסוד. מציינים את שם הסוד בשדה metadata > name בקובץ הזה.

List Secrets

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

gcloud beta composer environments user-workloads-secrets list \
  --environment ENVIRONMENT_NAME \
  --location LOCATION

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

  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.

קבלת פרטי הסוד

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

gcloud beta composer environments user-workloads-secrets describe \
  SECRET_NAME \
  --environment ENVIRONMENT_NAME \
  --location LOCATION

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

  • SECRET_NAME: השם של הסוד, כפי שהוגדר בשדה metadata > name בקובץ ה-YAML עם ההגדרה של הסוד.
  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.

מחיקת סוד

כדי למחוק סוד, מריצים את הפקודה הבאה:

gcloud beta composer environments user-workloads-secrets delete \
  SECRET_NAME \
  --environment ENVIRONMENT_NAME \
  --location LOCATION
  • SECRET_NAME: השם של הסוד, כפי שהוגדר בשדה metadata > name בקובץ ה-YAML עם ההגדרה של הסוד.
  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.

API

יצירת סוד

  1. יוצרים בקשת API של environments.userWorkloadsSecrets.create.

  2. בבקשה הזו:

    1. בגוף הבקשה, בשדה name, מציינים את ה-URI של הסוד החדש.
    2. בגוף הבקשה, בשדה data, מציינים מפתחות וערכים מקודדים בפורמט Base64 של הסוד.
    3. (אופציונלי, רק ב-Cloud Composer API בגרסת בטא) בגוף הבקשה, בשדה type, מציינים את סוג הסוד. אלה הסוגים של סודות שנתמכים:

      • (ברירת מחדל) Opaque
      • kubernetes.io/dockerconfigjson

דוגמה:

// POST https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsSecrets

{
  "name": "projects/example-project/locations/us-central1/environments/example-environment/userWorkloadsSecrets/example-secret",
  "data": {
    "example": "ZXhhbXBsZV92YWx1ZSAtbgo="
  }

}

דוגמה לסוג kubernetes.io/dockerconfigjson:

// POST https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsSecrets

{
  "name": "projects/example-project/locations/us-central1/environments/example-environment/userWorkloadsSecrets/example-secret",
  "data": {
    ".dockerconfigjson": "eyJhdXRocyI6eyJodHRwczovL2luZGV4LmRvY2tlci5pby92MS8iOnsidXNlcm5hbWUiOiJleGFtcGxlX3VzZXIiLCJwYXNzd29yZCI6ImV4YW1wbGVfcGFzc3dvcmQiLCJlbWFpbCI6ImVtYWlsQGV4YW1wbGUuY29tIiwiYXV0aCI6ImRYTmxjanB0ZVhCaGMzTT0ifX19"
  },
  "type": "kubernetes.io/dockerconfigjson"
}

עדכון סוד

  1. יוצרים בקשת API של environments.userWorkloadsSecrets.update.

  2. בבקשה הזו:

    1. בגוף הבקשה, בשדה name, מציינים את ה-URI של הסוד.
    2. בגוף הבקשה, בשדה data, מציינים מפתחות וערכים מקודדים בפורמט Base64 של הסוד. הערכים יוחלפו.

דוגמה:

// PUT https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsSecrets/example-secret

{
  "name": "projects/example-project/locations/us-central1/environments/example-environment/userWorkloadsSecrets/example-secret",
  "data": {
    "example": "ZXhhbXBsZV92YWx1ZSAtbgo=",
    "another-example": "YW5vdGhlcl9leGFtcGxlX3ZhbHVlIC1uCg=="
  }
}

List Secrets

יוצרים בקשת API של environments.userWorkloadsSecrets.list. ערכי מפתח בפלט יוחלפו בכוכביות. אפשר להשתמש בחלוקה לעמודים עם הבקשה הזו. לפרטים נוספים, אפשר לעיין בהפניה של הבקשה.

דוגמה:

// GET https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsSecrets

קבלת פרטי הסוד

יוצרים בקשת API של environments.userWorkloadsSecrets.get. ערכי מפתח בפלט יוחלפו בכוכביות.

דוגמה:

// GET https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsSecrets/example-secret

מחיקת סוד

יוצרים בקשת API של environments.userWorkloadsSecrets.delete.

דוגמה:

// DELETE https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsSecrets/example-secret

Terraform

המשאב google_composer_user_workloads_secret מגדיר סוד של Kubernetes, עם מפתחות וערכים שמוגדרים בבלוק data.

resource "google_composer_user_workloads_secret" "example_secret" {
  provider = google-beta
  environment = google_composer_environment.ENVIRONMENT_RESOURCE_NAME.name
  name = "SECRET_NAME"
  region = "LOCATION"

  data = {
    KEY_NAME: "KEY_VALUE"
  }
}
  • ENVIRONMENT_RESOURCE_NAME: השם של משאב הסביבה, שמכיל את ההגדרה של הסביבה ב-Terraform. שם הסביבה בפועל מצוין גם במשאב הזה.
  • LOCATION: האזור שבו נמצאת הסביבה.
  • SECRET_NAME: השם של הסוד.
  • KEY_NAME: מפתח אחד או יותר של הסוד הזה.
  • KEY_VALUE: ערך המפתח בקידוד base64. אפשר להשתמש בפונקציה base64encode כדי לקודד את הערך (ראו את הדוגמה).

שתי הדוגמאות הבאות של Kubernetes Secrets משמשות בדוגמאות בהמשך המדריך הזה.

resource "google_composer_user_workloads_secret" "example_secret" {
  provider = google-beta

  name = "airflow-secrets"

  environment = google_composer_environment.example_environment.name
  region = "us-central1"

  data = {
    sql_alchemy_conn: base64encode("postgresql+psycopg2://root:example-password@127.0.0.1:3306/example-db")
  }
}

דוגמה נוספת שמראה איך לכלול קבצים. אפשר להשתמש בפונקציה file כדי לקרוא את תוכן הקובץ כמחרוזת, ואז לקודד אותו ב-Base64:

resource "google_composer_user_workloads_secret" "service_account_secret" {
  provider = google-beta

  name = "service-account"

  environment = google_composer_environment.example_environment.name
  region = "us-central1"

  data = {
    "service-account.json": base64encode(file("./key.json"))
  }
}

שימוש ב-Kubernetes Secrets ב-DAG

בדוגמה הזו מוצגות שתי דרכים להשתמש ב-Kubernetes Secrets: כמשתנה סביבתי וכנפח שמוטמע על ידי ה-Pod.

הסוד הראשון, airflow-secrets, מוגדר כמשתנה סביבה של Kubernetes בשם SQL_CONN (בניגוד למשתנה סביבה של Airflow או Managed Airflow).

הסוד השני, service-account, מטמיע את service-account.json, קובץ עם טוקן של חשבון שירות, ב-/var/secrets/google.

כך נראים אובייקטים מסוג Secret:

secret_env = Secret(
    # Expose the secret as environment variable.
    deploy_type="env",
    # The name of the environment variable, since deploy_type is `env` rather
    # than `volume`.
    deploy_target="SQL_CONN",
    # Name of the Kubernetes Secret
    secret="airflow-secrets",
    # Key of a secret stored in this Secret object
    key="sql_alchemy_conn",
)
secret_volume = Secret(
    deploy_type="volume",
    # Path where we mount the secret as volume
    deploy_target="/var/secrets/google",
    # Name of Kubernetes Secret
    secret="service-account",
    # Key in the form of service account file name
    key="service-account.json",
)

השם של סוד Kubernetes הראשון מוגדר במשתנה secret_env. הסוד הזה נקרא airflow-secrets. הפרמטר deploy_type מציין שהמשתנה צריך להיות חשוף כמשתנה סביבה. השם של משתנה הסביבה הוא SQL_CONN, כפי שצוין בפרמטר deploy_target. לבסוף, הערך של משתנה הסביבה SQL_CONN מוגדר לערך של המפתח sql_alchemy_conn.

השם של הסוד השני ב-Kubernetes מוגדר במשתנה secret_volume. הסוד הזה נקרא service-account. הוא נחשף כנפח, כפי שמצוין בפרמטר deploy_type. הנתיב של הקובץ שצריך לטעון, deploy_target, הוא /var/secrets/google. לבסוף, הערך key של הסוד שמאוחסן ב-deploy_target הוא service-account.json.

כך נראית הגדרת האופרטור:

kubernetes_secret_vars_ex = KubernetesPodOperator(
    task_id="ex-kube-secrets",
    name="ex-kube-secrets",
    namespace="composer-user-workloads",
    image="marketplace.gcr.io/google/ubuntu2204",
    startup_timeout_seconds=300,
    # The secrets to pass to Pod, the Pod will fail to create if the
    # secrets you specify in a Secret object do not exist in Kubernetes.
    secrets=[secret_env, secret_volume],
    # Entrypoint of the container, if not specified the Docker container's
    # entrypoint is used. The cmds parameter is templated.
    cmds=["echo"],
    # env_vars allows you to specify environment variables for your
    # container to use. The env_vars parameter is templated.
    env_vars={
        "EXAMPLE_VAR": "/example/value",
        "GOOGLE_APPLICATION_CREDENTIALS": "/var/secrets/google/service-account.json",
    },
    # Specifies path to kubernetes config. The config_file is templated.
    config_file="/home/airflow/composer_kube_config",
    # Identifier of connection that should be used
    kubernetes_conn_id="kubernetes_default",
)

ניהול של ConfigMaps ב-Kubernetes

gcloud

יצירת ConfigMap

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

gcloud beta composer environments user-workloads-config-maps create \
  --environment ENVIRONMENT_NAME \
  --location LOCATION \
  --config-map-file-path CONFIG_MAP_FILE

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

  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.
  • CONFIG_MAP_FILE: נתיב לקובץ YAML מקומי שמכיל את התצורה של ConfigMap.

דוגמה:

gcloud beta composer environments user-workloads-config-maps create \
  --environment example-environment \
  --location us-central1 \
  --config-map-file-path ./configs/example-configmap.yaml

עדכון של ConfigMap

כדי לעדכן ConfigMap, מריצים את הפקודה הבאה. השם של ConfigMap יילקח מקובץ ה-YAML שצוין, והתוכן של ConfigMap יוחלף.

gcloud beta composer environments user-workloads-config-maps update \
  --environment ENVIRONMENT_NAME \
  --location LOCATION \
  --config-map-file-path CONFIG_MAP_FILE

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

  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.
  • CONFIG_MAP_FILE: נתיב לקובץ YAML מקומי שמכיל את התצורה של ConfigMap. מציינים את השם של ConfigMap בשדה metadata > name בקובץ הזה.

רשימת ConfigMaps

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

gcloud beta composer environments user-workloads-config-maps list \
  --environment ENVIRONMENT_NAME \
  --location LOCATION

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

  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.

קבלת פרטים של ConfigMap

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

gcloud beta composer environments user-workloads-config-maps describe \
  CONFIG_MAP_NAME \
  --environment ENVIRONMENT_NAME \
  --location LOCATION

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

  • CONFIG_MAP_NAME: השם של ConfigMap, כפי שהוגדר בשדה name > metadata בקובץ ה-YAML עם ההגדרה של ConfigMap.
  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.

מחיקת ConfigMap

כדי למחוק ConfigMap, מריצים את הפקודה הבאה:

gcloud beta composer environments user-workloads-config-maps delete \
  CONFIG_MAP_NAME \
  --environment ENVIRONMENT_NAME \
  --location LOCATION
  • CONFIG_MAP_NAME: השם של ConfigMap, כפי שהוגדר בשדה name > metadata בקובץ ה-YAML עם ההגדרה של ConfigMap.
  • ENVIRONMENT_NAME: השם של הסביבה.
  • LOCATION: האזור שבו נמצאת הסביבה.

API

יצירת ConfigMap

  1. יוצרים בקשת API של environments.userWorkloadsConfigMaps.create.

  2. בבקשה הזו:

    1. בגוף הבקשה, בשדה name, מציינים את ה-URI של ה-ConfigMap החדש.
    2. בגוף הבקשה, בשדה data, מציינים מפתחות וערכים עבור ConfigMap.

דוגמה:

// POST https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsConfigMaps

{
  "name": "projects/example-project/locations/us-central1/environments/example-environment/userWorkloadsConfigMaps/example-configmap",
  "data": {
    "example_key": "example_value"
  }
}

עדכון של ConfigMap

  1. יוצרים בקשת API של environments.userWorkloadsConfigMaps.update.

  2. בבקשה הזו:

    1. בגוף הבקשה, בשדה name, מציינים את ה-URI של ConfigMap.
    2. בגוף הבקשה, בשדה data, מציינים מפתחות וערכים עבור ConfigMap. הערכים יוחלפו.

דוגמה:

// PUT https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsConfigMaps/example-configmap

{
  "name": "projects/example-project/locations/us-central1/environments/example-environment/userWorkloadsConfigMaps/example-configmap",
  "data": {
    "example_key": "example_value",
    "another_key": "another_value"
  }
}

רשימת ConfigMaps

יוצרים בקשת API של environments.userWorkloadsConfigMaps.list. ערכי המפתח בפלט יוצגו כמו שהם. אפשר להשתמש בעימוד בבקשה הזו. פרטים נוספים זמינים בהפניה לבקשה.

דוגמה:

// GET https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsConfigMaps

קבלת פרטים של ConfigMap

יוצרים בקשת API של environments.userWorkloadsConfigMaps.get. ערכים מרכזיים בפלט יוצגו כמו שהם.

דוגמה:

// GET https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsConfigMaps/example-configmap

מחיקת ConfigMap

יוצרים בקשת API של environments.userWorkloadsConfigMaps.delete.

דוגמה:

// DELETE https://composer.googleapis.com/v1beta1/projects/example-project/
// locations/us-central1/environments/example-environment/userWorkloadsConfigMaps/example-configmap

Terraform

המשאב google_composer_user_workloads_config_map מגדיר ConfigMap, עם מפתחות וערכים שמוגדרים בבלוק data.

resource "google_composer_user_workloads_config_map" "example_config_map" {
  provider = google-beta
  environment = google_composer_environment.ENVIRONMENT_RESOURCE_NAME.name
  name = "CONFIG_MAP_NAME"
  region = "LOCATION"

  data = {
    KEY_NAME: "KEY_VALUE"
  }
}
  • ENVIRONMENT_RESOURCE_NAME: השם של משאב הסביבה, שמכיל את ההגדרה של הסביבה ב-Terraform. שם הסביבה בפועל מצוין גם במשאב הזה.
  • LOCATION: האזור שבו נמצאת הסביבה.
  • CONFIG_MAP_NAME: השם של ConfigMap.
  • KEY_NAME: מפתח אחד או יותר עבור ה-ConfigMap הזה.
  • KEY_VALUE: הערך של המפתח.

דוגמה:

resource "google_composer_user_workloads_config_map" "example_config_map" {
  provider = google-beta

  name = "example-config-map"

  environment = google_composer_environment.example_environment.name
  region = "us-central1"

  data = {
    "example_key": "example_value"
  }
}

שימוש ב-ConfigMaps ב-DAGs

בדוגמה הזו אפשר לראות איך משתמשים ב-ConfigMaps ב-DAGs.

בדוגמה הבאה, ConfigMap מועבר בפרמטר configmaps. כל המפתחות של ConfigMap זמינים כמשתני סביבה:

import datetime

from airflow import models
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator

with models.DAG(
    dag_id="composer_kubernetes_pod_configmap",
    schedule=None,
    start_date=datetime.datetime(2024, 1, 1),
) as dag:

  KubernetesPodOperator(
    task_id='kpo_configmap_env_vars',
    image='busybox:1.28',
    cmds=['sh'],
    arguments=[
        '-c',
        'echo "Value: $example_key"',
    ],
    configmaps=["example-configmap"],
    config_file="/home/airflow/composer_kube_config",
  )

בדוגמה הבאה אפשר לראות איך לטעון ConfigMap כנפח:

import datetime

from airflow import models
from kubernetes.client import models as k8s
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator

volume_mount = k8s.V1VolumeMount(name='confmap-example',
  mount_path='/config',
  sub_path=None,
  read_only=False)

volume = k8s.V1Volume(name='confmap-example',
  config_map=k8s.V1ConfigMapVolumeSource(name='example-configmap'))

with models.DAG(
    dag_id="composer_kubernetes_pod_configmap",
    schedule=None,
    start_date=datetime.datetime(2024, 1, 1),
) as dag:

  KubernetesPodOperator(
    task_id='kpo_configmap_volume_mount',
    image='busybox:1.28',
    cmds=['sh'],
    arguments=[
        '-c',
        'ls /config'
    ],
    volumes=[volume],
    volume_mounts=[volume_mount],
    configmaps=["example-configmap"],
    config_file="/home/airflow/composer_kube_config",
  )

מידע על ספק CNCF Kubernetes

האופרטור KubernetesPodOperator מיושם בספק apache-airflow-providers-cncf-kubernetes.

הערות מפורטות על הגרסה של ספק CNCF Kubernetes זמינות באתר של ספק CNCF Kubernetes.

משאבים נדרשים

ב-Managed Airflow (דור 3) יש תמיכה בערכים הבאים של דרישות משאבים. דוגמה לשימוש בדרישות משאבים מופיעה במאמר הגדרות נוספות.

משאב מינימום מקסימום שלב
CPU 0.25 32 ערכי השלבים: 0.25,‏ 0.5,‏ 1,‏ 2,‏ 4,‏ 6,‏ 8,‏ 10 וכן הלאה עד 32. הערכים המבוקשים מעוגלים כלפי מעלה לערך השלב הנתמך הקרוב ביותר (לדוגמה, 5 מעוגל ל-6).
זיכרון ‫2G (GB) ‫128G (GB) ערכי השלבים: 2,‏ 3,‏ 4,‏ 5 ועד 128. הערכים המבוקשים מעוגלים כלפי מעלה לערך השלב הנתמך הקרוב ביותר (לדוגמה, 3.5G מעוגל ל-4G).
אחסון - ‫100G (GB) כל ערך. אם מבקשים יותר מ-100GB, מקבלים רק 100GB.

מידע נוסף על יחידות משאבים ב-Kubernetes זמין במאמר יחידות משאבים ב-Kubernetes.

פתרון בעיות

בקטע הזה מפורטות הצעות לפתרון בעיות נפוצות שקשורות ל-KubernetesPodOperator:

צפייה ביומנים

כשמנסים לפתור בעיות, אפשר לבדוק את היומנים בסדר הבא:

  1. יומני משימות של Airflow:

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

      מעבר אל Environments

    2. ברשימת הסביבות, לוחצים על שם הסביבה. הדף Environment details ייפתח.

    3. עוברים לכרטיסייה DAGs.

    4. לוחצים על השם של ה-DAG ואז על ההרצה של ה-DAG כדי לראות את הפרטים והיומנים.

  2. יומנים של מתזמן Airflow:

    1. עוברים לדף פרטי הסביבה.

    2. עוברים לכרטיסייה יומנים.

    3. בודקים את היומנים של מתזמן Airflow.

  3. יומנים של עומסי עבודה של משתמשים:

    1. עוברים לדף פרטי הסביבה.

    2. עוברים לכרטיסייה מעקב.

    3. בוחרים באפשרות עומסי עבודה של משתמשים.

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

קודי החזרה שאינם אפס

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

דפוס נפוץ הוא להריץ סקריפט של מעטפת כנקודת הכניסה של הקונטיינר כדי לקבץ כמה פעולות בתוך הקונטיינר.

אם אתם כותבים סקריפט כזה, מומלץ לכלול את הפקודה set -e בראש הסקריפט, כדי שפקודות שנכשלו בסקריפט יסיימו את הסקריפט ויעבירו את הכשל למופע המשימה של Airflow.

הזמן הקצוב לתפוגה של Pod

זמן הקצוב לתפוגה (timeout) שמוגדר כברירת מחדל ב-KubernetesPodOperator הוא 120 שניות, מה שיכול לגרום לזמני קצוב לתפוגה לפני הורדה של תמונות גדולות יותר. אפשר להגדיל את הזמן הקצוב לתפוגה באמצעות שינוי הפרמטר startup_timeout_seconds כשיוצרים את KubernetesPodOperator.

כשפוד מגיע לזמן קצוב לתפוגה, היומן הספציפי למשימה זמין בממשק המשתמש של Airflow. לדוגמה:

Executing <Task(KubernetesPodOperator): ex-all-configs> on 2018-07-23 19:06:58.133811
Running: ['bash', '-c', u'airflow run kubernetes-pod-example ex-all-configs 2018-07-23T19:06:58.133811 --job_id 726 --raw -sd DAGS_FOLDER/kubernetes_pod_operator_sample.py']
Event: pod-name-9a8e9d06 had an event of type Pending
...
...
Event: pod-name-9a8e9d06 had an event of type Pending
Traceback (most recent call last):
  File "/usr/local/bin/airflow", line 27, in <module>
    args.func(args)
  File "/usr/local/lib/python2.7/site-packages/airflow/bin/cli.py", line 392, in run
    pool=args.pool,
  File "/usr/local/lib/python2.7/site-packages/airflow/utils/db.py", line 50, in wrapper
    result = func(*args, **kwargs)
  File "/usr/local/lib/python2.7/site-packages/airflow/models.py", line 1492, in _run_raw_task
    result = task_copy.execute(context=context)
  File "/usr/local/lib/python2.7/site-packages/airflow/contrib/operators/kubernetes_pod_operator.py", line 123, in execute
    raise AirflowException('Pod Launching failed: {error}'.format(error=ex))
airflow.exceptions.AirflowException: Pod Launching failed: Pod took too long to start

יכול להיות שפסק זמן של Pod יתרחש גם אם ל-Managed Airflow חסרות הרשאות ה-IAM הנדרשות לביצוע המשימה. כדי לוודא זאת, אפשר להשתמש בלוחות הבקרה של GKE כדי לבדוק את היומנים של עומס העבודה הספציפי, או להשתמש ב-Cloud Logging.

משימות KubernetesPodOperator נכשלות כשמבצעים מספר גדול של משימות

כשבסביבה שלכם מופעל מספר גדול של משימות KubernetesPodOperator או KubernetesExecutor בו-זמנית, Managed Airflow (דור 3) לא מקבל משימות חדשות עד שחלק מהמשימות הקיימות מסתיימות.

מידע נוסף על פתרון הבעיה הזו זמין במאמר פתרון בעיות במשימות של KubernetesExecutor.

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