העברת נתונים באמצעות Google Transfer Operators

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:

  1. בממשק המשתמש של Airflow, עוברים אל Admin > Connections.
  2. יוצרים חיבור חדש.
  3. בוחרים באפשרות Amazon S3 כסוג החיבור.
  4. בדוגמה הבאה נעשה שימוש בחיבור בשם aws_s3. אפשר להשתמש בשם הזה או בכל שם אחר לחיבור.
  5. מציינים את פרמטרי החיבור כפי שמתואר במסמכי התיעוד של 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:

  1. בממשק המשתמש של Airflow, עוברים אל Admin > Connections.
  2. יוצרים חיבור חדש.
  3. בוחרים באפשרות Azure FileShare כסוג החיבור.
  4. בדוגמה הבאה נעשה שימוש בחיבור בשם azure_fileshare. אפשר להשתמש בשם הזה או בכל שם אחר לחיבור.
  5. מציינים את פרמטרי החיבור כמו שמתואר במסמכי 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

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