Managed Airflow (דור 3) | Managed Airflow (דור 2) | Managed Airflow (דור 1 מדור קודם)
בדף הזה מוסבר איך להשתמש בפונקציות של Cloud Run כדי להפעיל DAG של Managed Service for Apache Airflow בתגובה לאירועים.
Apache Airflow מיועד להרצת DAGs בלוח זמנים קבוע, אבל אפשר גם להפעיל DAGs בתגובה לאירועים. אחת הדרכים לעשות את זה היא להשתמש בפונקציות Cloud Run כדי להפעיל DAGs מנוהלים של Airflow כשמתרחש אירוע ספציפי.
אפשר גם:
- הפעלת DAG באמצעות API בארכיטקטורת REST של Airflow בלבד.
- יצירת פונקציה שמפעילה DAG כשדוחפים הודעה לנושא Pub/Sub.
בדוגמה שבמדריך הזה מוצגת פונקציה שמפעילה DAG בתגובה לאירוע:
- מגדירים טריגרים לפונקציה ב-Cloud Run Functions.
- כשהפונקציה מופעלת, היא שולחת בקשה להפעלת DAG דרך API בארכיטקטורת REST של Airflow בסביבת Managed Airflow. הבקשה מכילה את המזהה והסוג של האירוע, ואת מטען הייעודי (payload) של האירוע.
- מערכת Airflow מעבדת את הבקשה ומריצה את ה-DAG שצוין בבקשה. ה-DAG מוציא את הנתונים שהועברו אליו מהפונקציה.
לפני שמתחילים
בקטע הזה מפורטים שלבי ההכנה.
בדיקת הגדרות הרשת בסביבה
הפתרון הזה לא פועל בהגדרות של כתובות IP פרטיות ו-VPC Service Controls, כי אי אפשר להגדיר קישוריות מפונקציות של Cloud Run לשרת האינטרנט של Airflow בהגדרות האלה.
הפעלת ממשקי API בפרויקט
המסוף
מפעילים את ממשקי ה-API של Managed Airflow ושל פונקציות Cloud Run.
תפקידים שנדרשים להפעלת ממשקי API
כדי להפעיל ממשקי API, נדרשת ההרשאה serviceusage.services.enable. אם יצרתם את הפרויקט, סביר להניח שכבר יש לכם את ההרשאה הזו דרך התפקיד 'בעלים' (roles/owner). אחרת, תוכלו לקבל את ההרשאה הזו דרך התפקיד 'אדמין בממשק Service Usage' (roles/serviceusage.serviceUsageAdmin). איך מקצים תפקידים
gcloud
מפעילים את ממשקי ה-API של פונקציות Managed Airflow ו-Cloud Run:
תפקידים שנדרשים להפעלת ממשקי API
כדי להפעיל ממשקי API, נדרשת ההרשאה serviceusage.services.enable. אם יצרתם את הפרויקט, סביר להניח שכבר יש לכם את ההרשאה הזו דרך התפקיד 'בעלים' (roles/owner). אחרת, תוכלו לקבל את ההרשאה הזו דרך התפקיד 'אדמין בממשק Service Usage' (roles/serviceusage.serviceUsageAdmin). איך מקצים תפקידים
gcloud services enable cloudfunctions.googleapis.comcomposer.googleapis.com
הפעלת Airflow API בארכיטקטורת REST
בהתאם לגרסת Airflow:
- ב-Airflow 2, ממשק ה-API היציב ל-REST כבר מופעל כברירת מחדל. אם ה-API היציב מושבת בסביבה שלכם, צריך להפעיל את ה-API היציב ל-REST.
- ב-Airflow 1, מפעילים את ה-API הניסיוני ל-REST.
התרת קריאות ל-API בארכיטקטורת REST של Airflow באמצעות בקרת גישה לרשת של שרת האינטרנט
פונקציות Cloud Run יכולות לפנות אל Airflow API בארכיטקטורת REST באמצעות כתובת IPv4 או IPv6.
אם אתם לא בטוחים מה יהיה טווח כתובות ה-IP של הקריאה, כדאי להשתמש באפשרות ברירת המחדל להגדרה בWebserver Access Control, שהיא All IP addresses have access (default), כדי שלא תחסמו בטעות את הפונקציות של Cloud Run. תמיד אפשר להגדיר גישה לרשת של שרת האינטרנט מאוחר יותר.
איך מוצאים את כתובת ה-URL של שרת האינטרנט של Airflow
בדוגמה הזו מתבצעות בקשות ל-API בארכיטקטורת REST לנקודת הקצה של שרת האינטרנט של Airflow.
משתמשים בחלק של כתובת ה-URL של ממשק האינטרנט של Airflow לפני .appspot.com בקוד של Cloud Function.
המסוף
נכנסים לדף Environments במסוף Google Cloud .
לוחצים על שם הסביבה.
בדף Environment details, עוברים לכרטיסייה Environment configuration.
כתובת ה-URL של שרת האינטרנט של Airflow מופיעה בפריט Airflow web UI.
gcloud
מריצים את הפקודה הבאה:
gcloud composer environments describe ENVIRONMENT_NAME \
--location LOCATION \
--format='value(config.airflowUri)'
מחליפים את:
-
ENVIRONMENT_NAMEבשם הסביבה. -
LOCATIONעם האזור שבו הסביבה ממוקמת.
קבלת מזהה הלקוח של שרת ה-IAM proxy
כדי לשלוח בקשה לנקודת קצה ל-API בארכיטקטורת REST של Airflow, הפונקציה צריכה את מזהה הלקוח של שרת proxy לניהול זהויות והרשאות גישה שמגן על שרת האינטרנט של Airflow.
ב-Managed Airflow, המידע הזה לא זמין באופן ישיר. במקום זאת, שולחים בקשה לא מאומתת לשרת האינטרנט של Airflow ומתעדים את מזהה הלקוח מכתובת ה-URL להפניה אוטומטית:
cURL
curl -v AIRFLOW_URL 2>&1 >/dev/null | grep -o "client_id\=[A-Za-z0-9-]*\.apps\.googleusercontent\.com"
מחליפים את AIRFLOW_URL בכתובת ה-URL של ממשק האינטרנט של Airflow.
בפלט, מחפשים את המחרוזת שאחרי client_id. לדוגמה:
client_id=836436932391-16q2c5f5dcsfnel77va9bvf4j280t35c.apps.googleusercontent.com
Python
שומרים את הקוד הבא בקובץ בשם get_client_id.py. ממלאים את הערכים של project_id, location ו-composer_environment, ואז מריצים את הקוד ב-Cloud Shell או בסביבה המקומית.
העלאת DAG לסביבה
העלאת DAG לסביבה הפלט של ה-DAG הבא הוא הגדרת ההרצה של ה-DAG שהתקבלה. אפשר להפעיל את ה-DAG הזה מפונקציה שיוצרים בהמשך המדריך הזה.
import datetime
import airflow
from airflow.operators.bash_operator import BashOperator
with airflow.DAG(
'composer_sample_trigger_response_dag',
start_date=datetime.datetime(2026, 1, 1),
# Not scheduled, trigger only
schedule=None) as dag:
# Print the dag_run's configuration, which includes information about the
# Cloud Storage object change.
print_gcs_info = BashOperator(
task_id='print_gcs_info', bash_command='echo {{ dag_run.conf }}}}')
פריסת פונקציה שמפעילה את ה-DAG
אפשר לפרוס פונקציה באמצעות השפה המועדפת שנתמכת על ידי פונקציות Cloud Run או Cloud Run. במדריך הזה נסביר איך להטמיע Cloud Function ב-Python וב-Java.
ציון פרמטרים להגדרת פונקציה
טריגר: בוחרים טריגר אחד או יותר של Eventarc לפונקציה.
מידע נוסף על יצירת טריגרים זמין במאמר יצירת טריגרים באמצעות Eventarc. לדוגמה, אפשר להפעיל פונקציות מ-Cloud Storage באמצעות Eventarc.
חשבון שירות: לחשבון השירות שאתם מציינים לטריגר צריכות להיות הרשאות מספיקות להפעלת DAG בסביבות Managed Airflow.
מומלץ לפעול לפי העיקרון של הרשאות מינימליות ולהעניק לו רק את התפקיד Composer User (
composer.user). מידע נוסף על הגדרת הרשאות זמין במאמר תפקידים והרשאות ליעדים ב-Cloud Run.Function entry point:
(Python) כשמוסיפים קוד לדוגמה הזו, בוחרים את זמן הריצה Python 3.10 או גרסה מאוחרת יותר ומציינים את
trigger_dagכנקודת הכניסה.
הוספת דרישות
מציינים את יחסי התלות בקובץ requirements.txt:
הוספת קוד פונקציה
מוסיפים את הקוד הבא לקובץ main.py ומבצעים את ההחלפות הבאות:
מחליפים את הערך של המשתנה
client_idבערךclient_idשקיבלתם קודם.מחליפים את הערך של המשתנה
webserver_idבמזהה פרויקט הדייר, שהוא חלק מכתובת ה-URL של ממשק האינטרנט של Airflow לפני.appspot.com. כתובת ה-URL של ממשק האינטרנט של Airflow הופיעה קודם.מציינים את הגרסה של Airflow API בארכיטקטורת REST שבה משתמשים:
- אם אתם משתמשים ב-Airflow API בארכיטקטורת REST היציב, צריך להגדיר את המשתנה
USE_EXPERIMENTAL_APIלערךFalse. - אם אתם משתמשים ב-Airflow API בארכיטקטורת REST הניסיוני, לא צריך לבצע שינויים. המשתנה
USE_EXPERIMENTAL_APIכבר מוגדר לערךTrue.
- אם אתם משתמשים ב-Airflow API בארכיטקטורת REST היציב, צריך להגדיר את המשתנה
בדיקת הפונקציה
כדי לוודא שהפונקציה ו-DAG פועלים כמצופה:
- מחכים עד שהפונקציה תופעל.
- מפעילים את הפונקציה בהתאם לטריגר שצוין. אפשר גם להפעיל את הפונקציה באופן ידני על ידי בחירה בפעולה Test the function (בדיקת הפונקציה) עבורה ב- Google Cloud console.
- בודקים את דף ה-DAG בממשק האינטרנט של Airflow. ל-DAG צריך להיות הפעלה אחת פעילה או שכבר הסתיימה.
- בממשק המשתמש של Airflow, בודקים את יומני המשימות של ההרצה הזו. אפשר לראות שהנתונים שהתקבלו מהפונקציה מופיעים ביומנים של המשימה
print_gcs_info:
פלט לדוגמה:
[2021-04-04 18:25:44,778] {bash_operator.py:154} INFO - Output:
[2021-04-04 18:25:44,781] {bash_operator.py:158} INFO - Triggered from GCF:
{bucket: example-storage-for-gcf-triggers, contentType: text/plain,
crc32c: dldNmg==, etag: COW+26Sb5e8CEAE=, generation: 1617560727904101,
... }
[2021-04-04 18:25:44,781] {bash_operator.py:162} INFO - Command exited with
return code 0h
המאמרים הבאים
- גישה לממשק המשתמש של Airflow
- גישה ל-Airflow API בארכיטקטורת REST
- כתיבת תרשימי DAG
- כתיבה של פונקציות Cloud Run
- טריגרים של Cloud Storage