Apache Ecosystem

حرفه‌ای شدن در Apache Airflow: ساخت پایپ‌لاین‌های داده مقیاس‌پذیر از صفر

با انفجار حجم داده‌ها و توزیع‌یابی فزاینده‌ی زیرساخت‌ها، نیاز به هماهنگ‌سازی قابل اعتماد، قابل مشاهده و قابل توسعه‌ی جریان‌های کاری هرگز این‌قدر بزرگ نبوده است. Apache Airflow به عنوان استاندارد صنعتی در این زمینه ظهور کرده است. با بیان جریان‌های کاری به صورت گراف‌های جهت‌دار غیرسیکلیک (DAG) با کد پایتون، Airflow به مهندسان داده امکان می‌دهد از اسکریپت‌های شکننده shell فاصله بگیرند و به پایپ‌لاین‌های داده‌ای تحت کنترل نسخه، آزمایش‌شده و قابل نگهداری برسند.

چه در حال اتوماسیون شغل‌های بچ ساده‌ی شبانه باشید و چه هماهنگ‌سازی میکروسرویس‌های پیچیده‌ی بلادرنگ، درک اجزای اصلی Airflow—DAGها، اپراتورها، سنسورها و استراتژی‌های اجرا—برای هر استک داده‌ای مدرن ضروری است. در این راهنما، معماری را تجزیه و تحلیل می‌کنیم، نمونه‌های کد عملی را بررسی می‌کنیم و ملاحظات حیاتی برای استقرار در محیط تولید را بحث می‌کنیم.

انتزاع اصلی: DAGها و زمان‌بندی

در قلب Airflow، DAG (گراف جهت‌دار غیرسیکلیک) قرار دارد. یک DAG مجموعه‌ای از وظایف با وابستگی‌هاست که به صورت گره‌ها (وظایف) و یال‌ها (وابستگی‌ها) نمایش داده می‌شوند. برخلاف سیستم‌های زمان‌بندی سنتی که به فایل‌های پیکربندی ایستا متکی هستند، Airflow DAGها را در پایتون تعریف می‌کند. این امکان را فراهم می‌کند که وظایف به صورت پویا تولید شوند، منطق شرطی و پارامترسازی مستقیماً درون کد شما انجام شود.

وقتی یک DAG را در رابط کاربری Airflow استقرار می‌دهید، موتور زمان‌بندی پوشه‌ی ~/.airflow/dags را اسکن می‌کند. فایل پایتون را تجزیه می‌کند، زمان‌بندی را تعیین می‌کند (مثلاً به سبک cron) و برای هر بازه‌ی زمان‌بندی شده "اجراهای DAG" ایجاد می‌کند. درک پارامترهای start_date و schedule_interval حیاتی است، زیرا آن‌ها زمینه‌ی زمانی برای پایپ‌لاین داده شما را تعریف می‌کنند.

اجزای سازنده: اپراتورها، سنسورها و هوک‌ها

وظایف در Airflow توسط اپراتورها اجرا می‌شوند. یک اپراتور یک کلاس پایتون است که یک اقدام یا مرحله‌ی واحد در پایپ‌لاین را تعریف می‌کند. Airflow اکوسیستم گسترده‌ای از اپراتورهای از پیش ساخته‌شده برای سرویس‌هایی مانند AWS، GCP، DBT، Kubernetes و SQL ارائه می‌دهد. با این حال، وقتی منطق خاصی مورد نیاز باشد، می‌توانید اپراتورهای سفارشی بنویسید یا از PythonOperator استفاده کنید.

اغلب اوقات، جریان‌های کاری به رویدادهای خارجی وابسته هستند که ممکن است دقیقاً در زمان زمان‌بندی‌شده رخ ندهند. سنسورها با پل کردن (polling) یک سیستم خارجی (مانند یک سطل S3 یا یک پایگاه داده) تا زمانی که یک شرط برقرار شود، قبل از فعال‌سازی وظیفه‌ی بعدی، این مشکل را حل می‌کنند. این جداسازی، یکپارچگی داده را تضمین می‌کند و از شرایط مسابقه (race conditions) در محیط‌های توزیع‌شده جلوگیری می‌کند.

هوک‌ها (Hooks) یک رابط سطح پایین به سرویس‌های خارجی ارائه می‌دهند، احراز هویت و فراخوانی‌های API را مدیریت می‌کنند و سپس توسط اپراتورها و سنسورها برای انتزاع سطح بالاتر پوشش داده می‌شوند.

نمونه عملی: یک پایپ‌لاین ETL

یک سناریوی رایج ETL را در نظر بگیرید: استخراج داده از یک پایگاه داده PostgreSQL، تبدیل آن در Pandas و بارگذاری آن در یک دریاچه داده S3.

from airflow import DAG
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.providers.apache.hive.operators.hive import HiveOperator
from airflow.operators.python import PythonOperator
from airflow.providers.aws.operators.s3 import S3Operator
from airflow.sensors.postgres import PostgresSensor
from datetime import datetime, timedelta

default_args = {
    'owner': 'data-team',
    'retries': 3,
    'retry_delay': timedelta(minutes=5)
}

def extract_custom_logic(**context):
    # Custom Python logic for transformation
    data = context['task_instance'].xcom_pull(task_ids='extract_sql')
    processed_data = transform(data) 
    return processed_data

def load_to_s3(processed_data):
    # Logic to upload to S3
    pass

with DAG(
    dag_id='etl_sales_pipeline',
    default_args=default_args,
    description='ETL pipeline for sales data',
    schedule_interval='@daily',
    start_date=datetime(2023, 1, 1),
    catchup=False,
    tags=['etl', 'sales'],
) as dag:

    wait_for_data = PostgresSensor(
        task_id='wait_for_new_data',
        sql="SELECT 1 FROM sales WHERE updated_at > now() - interval '1 day'",
        postgres_conn_id='prod_postgres',
        poke_interval=60,
    )

    extract = PostgresOperator(
        task_id='extract_sql',
        sql="SELECT * FROM sales WHERE updated_at > now() - interval '1 day'",
        postgres_conn_id='prod_postgres',
    )

    transform = PythonOperator(
        task_id='transform_data',
        python_callable=extract_custom_logic,
    )

    load = S3Operator(
        task_id='load_to_s3',
        s3_key='data/sales/etl_output.parquet',
        local_location='/tmp/sales_output.parquet',
    )

    wait_for_data >> extract >> transform >> load

در این مثال، ما یک پایپ‌لاین قوی با تلاش مجدد (retries) و یک سنسور برای اطمینان از در دسترس بودن داده تعریف می‌کنیم. xcom_pull به اپراتور پایتون اجازه می‌دهد داده را از وظیفه SQL قبلی دریافت کند و نشان می‌دهد که چگونه وضعیت بین مراحل منتقل می‌شود.

بهترین روش‌های استقرار در محیط تولید

انتقال Airflow به محیط تولید نیاز به چیزهای بیشتری دارد تا صرفاً نصب نرم‌افزار. بهترین روش‌های کلیدی عبارتند از:

  • جداسازی نگرانی‌ها: تعاریف DAGها را در Git نگه دارید، جدا از نمونه‌های worker Airflow.
  • انتخاب اجرایی‌کننده (Executor): برای مقیاس‌بندی افقی از اجرایی‌کننده Celery یا Kubernetes استفاده کنید و از اجرایی‌کننده Single Process برای بارهای کاری تولیدی خودداری کنید.
  • مدیریت اتصالات: رازها (secrets) را در متغیرهای محیطی یا یک مدیر راز (مانند AWS Secrets Manager) ذخیره کنید و از طریق connections به آن‌ها ارجاع دهید، به جای سخت‌کد کردن اعتبارنامه‌ها.
  • پایش: Airflow را با ابزارهای پایش مانند Prometheus، Grafana یا Datadog یکپارچه کنید تا مدت زمان وظایف، نرخ شکست و مصرف منابع را ردیابی کنید.
  • آزمایش: از چارچوب‌هایی مانند pytest برای آزمایش واحد منطق DAG و airflow dags test برای آزمایش‌های یکپارچه قبل از استقرار استفاده کنید.

نتیجه‌گیری

Apache Airflow یک چارچوب قدرتمند، انعطاف‌پذیر و قابل توسعه برای هماهنگ‌سازی جریان‌های کاری داده ارائه می‌دهد. با بهره‌گیری از سینتکس اعلایی DAG، اکوسیستم غنی اپراتورها و قابلیت‌های قوی زمان‌بندی، تیم‌های مهندسی داده می‌توانند پایپ‌لاین‌هایی بسازند که نه تنها کارآمد بلکه مقاوم و قابل نگهداری نیز هستند. با رشد زیرساخت داده شما، سرمایه‌گذاری در طراحی مناسب Airflow و تقویت آن برای محیط تولید، سودهای قابل توجهی در زمینه قابلیت اطمینان و کارایی عملیاتی خواهد داشت.

Share: