Managed Airflow (דור 3) | Managed Airflow (דור 2) | Managed Airflow (דור 1 מדור קודם)
בדף הזה מוסבר איך להעביר נתונים משירותים אחרים באמצעות Google Transfer Operators ב-DAG.
מידע על מפעילים של Google Transfer
Google Transfer Operators הם קבוצה של אופרטורים של Airflow שאפשר להשתמש בהם כדי לשלוף נתונים משירותים אחרים אלGoogle Cloud.
במדריך הזה מוצגים אופרטורים ל-Azure FileShare Storage ול-Amazon S3 שפועלים עם Cloud Storage. יש עוד הרבה ספקי העברה שעובדים עם שירותים בתוך Google Cloud ועם שירותים אחרים מלבדGoogle Cloud.
Amazon S3 ל-Cloud Storage
בקטע הזה מוסבר איך לסנכרן נתונים מ-Amazon S3 לקטגוריה של Cloud Storage.
התקנת חבילת הספק של אמזון
חבילת apache-airflow-providers-amazon מכילה את סוגי החיבור והפונקציונליות שפועלים באינטראקציה עם Amazon S3.
מתקינים את חבילת PyPI הזו בסביבה.
הגדרת חיבור ל-Amazon S3
חבילת הספק של אמזון מספקת סוג חיבור ל-Amazon S3. יוצרים חיבור מהסוג הזה. החיבור ל-Cloud Storage, שנקרא google_cloud_default, כבר מוגדר בסביבה שלכם.
כדי להגדיר חיבור ל-Amazon S3:
- בממשק המשתמש של Airflow, עוברים אל Admin > Connections.
- יוצרים חיבור חדש.
- בוחרים באפשרות
Amazon S3כסוג החיבור. - בדוגמה הבאה נעשה שימוש בחיבור בשם
aws_s3. אפשר להשתמש בשם הזה או בכל שם אחר לחיבור. - מציינים את פרמטרי החיבור כפי שמתואר במסמכי התיעוד של Airflow בנושא Amazon Web Services Connection. לדוגמה, כדי להגדיר חיבור עם מפתחות גישה ל-AWS, צריך ליצור מפתח גישה לחשבון ב-AWS, ואז לספק את מזהה מפתח הגישה ל-AWS בתור שם משתמש ואת מפתח הגישה הסודי ל-AWS בתור סיסמה לחיבור.
העברת נתונים מ-Amazon S3
אם רוצים לבצע פעולות על הנתונים המסונכרנים בשלב מאוחר יותר ב-DAG או במשימה אחרים, צריך למשוך אותם לתיקייה /data בקטגוריה של הסביבה. התיקייה הזו מסונכרנת עם עובדי Airflow אחרים, כך שמשימות ב-DAG יכולות לפעול בה.
ה-DAG הבא לדוגמה מבצע את הפעולות הבאות:
- מסנכרן את התוכן של הספרייה
/data-for-gcsמקטגוריית S3 עם התיקייה/data/from-s3/data-for-gcs/בקטגוריה של הסביבה. - ממתינים שתי דקות עד שהנתונים יסונכרנו עם כל העובדים של Airflow בסביבה.
- הפקודה
lsמציגה את רשימת הקבצים בספרייה הזו. מחליפים את המשימה הזו באופרטורים אחרים של Airflow שפועלים עם הנתונים שלכם.
import datetime
import airflow
from airflow.providers.google.cloud.transfers.s3_to_gcs import S3ToGCSOperator
from airflow.operators.bash_operator import BashOperator
with airflow.DAG(
'composer_sample_aws_to_gcs',
start_date=datetime.datetime(2022, 1, 1),
schedule=None,
) as dag:
transfer_dir_from_s3 = S3ToGCSOperator(
task_id='transfer_dir_from_s3',
aws_conn_id='aws_s3',
prefix='data-for-gcs',
bucket='example-s3-bucket-transfer-operators',
dest_gcs='gs://us-central1-example-environ-361f2312-bucket/data/from-s3/')
sleep_2min = BashOperator(
task_id='sleep_2min',
bash_command='sleep 2m')
print_dir_files = BashOperator(
task_id='print_dir_files',
bash_command='ls /home/airflow/gcs/data/from-s3/data-for-gcs/')
transfer_dir_from_s3 >> sleep_2min >> print_dir_files
מ-Azure FileShare אל Cloud Storage
בקטע הזה מוסבר איך לסנכרן נתונים מ-Azure FileShare לקטגוריה של Cloud Storage.
התקנת חבילת הספק של Microsoft Azure
חבילת apache-airflow-providers-microsoft-azure מכילה את סוגי החיבור והפונקציונליות שפועלים באינטראקציה עם Microsoft Azure.
מתקינים את חבילת PyPI הזו בסביבה.
הגדרת חיבור ל-Azure FileShare
חבילת הספק של Microsoft Azure מספקת סוג חיבור ל-Azure File Share. אתם יוצרים חיבור מהסוג הזה. החיבור ל-Cloud Storage בשם google_cloud_default כבר מוגדר בסביבה שלכם.
כדי להגדיר חיבור ל-Azure FileShare:
- בממשק המשתמש של Airflow, עוברים אל Admin > Connections.
- יוצרים חיבור חדש.
- בוחרים באפשרות
Azure FileShareכסוג החיבור. - בדוגמה הבאה נעשה שימוש בחיבור בשם
azure_fileshare. אפשר להשתמש בשם הזה או בכל שם אחר לחיבור. - מציינים את פרמטרי החיבור כמו שמתואר במסמכי Airflow בנושא Microsoft Azure File Share Connection. לדוגמה, אפשר לציין מחרוזת חיבור למפתח הגישה של חשבון האחסון.
העברת נתונים מ-Azure FileShare
אם רוצים לבצע פעולות על הנתונים המסונכרנים בשלב מאוחר יותר ב-DAG או במשימה אחרים, צריך למשוך אותם לתיקייה /data בקטגוריה של הסביבה. התיקייה הזו מסונכרנת עם עובדי Airflow אחרים, כך שמשימות ב-DAG יכולות לפעול בה.
ה-DAG הבא מבצע את הפעולות הבאות:
ה-DAG הבא לדוגמה מבצע את הפעולות הבאות:
- מסנכרן את התוכן של הספרייה
/data-for-gcsמ-Azure File Share לתיקייה/data/from-azureבדלי של הסביבה. - ממתינים שתי דקות עד שהנתונים יסונכרנו עם כל העובדים של Airflow בסביבה.
- הפקודה
lsמציגה את רשימת הקבצים בספרייה הזו. מחליפים את המשימה הזו באופרטורים אחרים של Airflow שפועלים עם הנתונים שלכם.
import datetime
import airflow
from airflow.providers.google.cloud.transfers.azure_fileshare_to_gcs import AzureFileShareToGCSOperator
from airflow.operators.bash_operator import BashOperator
with airflow.DAG(
'composer_sample_azure_to_gcs',
start_date=datetime.datetime(2022, 1, 1),
schedule=None,
) as dag:
transfer_dir_from_azure = AzureFileShareToGCSOperator(
task_id='transfer_dir_from_azure',
azure_fileshare_conn_id='azure_fileshare',
share_name='example-file-share',
directory_name='data-for-gcs',
dest_gcs='gs://us-central1-example-environ-361f2312-bucket/data/from-azure/')
sleep_2min = BashOperator(
task_id='sleep_2min',
bash_command='sleep 2m')
print_dir_files = BashOperator(
task_id='print_dir_files',
bash_command='ls /home/airflow/gcs/data/from-azure/')
transfer_dir_from_azure >> sleep_2min >> print_dir_files