Apache Ecosystem

تسلط بر Apache Airflow: از DAGها تا پایپ‌لاین‌های ETL در سطح تولید

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

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

در قلب هر نمونه Airflow، گراف جهت‌دار بدون دور (DAG) قرار دارد. یک DAG جریان کاری شما را به عنوان مجموعه‌ای از وظایف و وابستگی‌های آن‌ها تعریف می‌کند. جنبه «جهت‌دار» تضمین می‌کند که وابستگی‌های دایره‌ای وجود ندارد، در حالی که «بدون دور» بودن تضمین می‌کند که هر وظیفه از یک نقطه شروع قابل دسترسی است. زمان‌بندی در Airflow مبتنی بر زمان است. شما یک فاصله زمانی زمان‌بندی (با نحو شبیه به cron) تعریف می‌کنید و زمان‌بند تعیین می‌کند که هر اجرای DAG چه زمانی باید آغاز شود. این جداسازی زمان اجرا از زمان تعریف، امکان اجرای پایپ‌لاین‌های ایدمپوتنت (بدون اثر جانبی تکراری) را فراهم می‌کند. برای مثال، اگر یک پایپ‌لاین به دلیل یک خطای موقت شبکه شکست بخورد، اجرای مجدد همان بازه منطقی، در صورت طراحی صحیح، داده‌ها را تکرار نخواهد کرد.
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

default_args = {
    'owner': 'data_engineer',
    'depends_on_past': False,
    'start_date': datetime(2023, 1, 1),
    'retries': 3,
    'retry_delay': timedelta(minutes=5)
}

with DAG(
    'etl_daily_summary',
    default_args=default_args,
    description='Daily ETL pipeline for sales data',
    schedule_interval='@daily',
    catchup=False
) as dag:
    def extract_data(**kwargs):
        print("Extracting raw sales data...")

    extract_task = PythonOperator(
        task_id='extract',
        python_callable=extract_data
    )

عملگرها، سنسورها و وابستگی‌های وظیفه

وظایف درون یک DAG توسط اجزای مختلفی به نام عملگرها (Operators) اجرا می‌شوند. اگرچه PythonOperator انعطاف‌پذیرترین گزینه است، Airflow عملگرهای تخصصی برای SQL (مانند PostgresOperator و MySqlOperator)، ذخیره‌سازی ابری (مانند S3Hook) و کوبرنیتس (مانند KubernetesPodOperator) ارائه می‌دهد. برای وابستگی‌هایی که به شرایط خارجی به جای زمان ثابت یا تکمیل وظیفه قبلی وابسته هستند، از سنسورها استفاده می‌شود. یک سنسور منتظر می‌ماند تا شرط خاصی برآورده شود، مانند رسیدن یک فایل به یک سطل S3 یا به‌روزرسانی یک رکورد در پایگاه داده. این موضوع برای معماری‌های داده‌ای رویداد-محور حیاتی است.
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor

wait_for_data = S3KeySensor(
    task_id='wait_for_new_file',
    bucket_key='data/incoming/sales_*.csv',
    bucket_name='my-data-lake',
    timeout=3600,
    poke_interval=60
)

خودکارسازی ETL و بهترین شیوه‌ها

ساخت یک پایپ‌لاین ETL در Airflow نیازمند جداسازی دقیق نگرانی‌ها (Separation of Concerns) است. فازهای استخراج، تبدیل و بارگذاری (Extract, Transform, Load) باید اغلب به عنوان وظایف مجزا باشند تا امکان منطق تلاش مجدد و نظارت دقیق‌تر فراهم شود. از نوشتن اسکریپت‌های پیچیده پایتون که همه کارها را انجام می‌دهند، پرهیز کنید. در عوض، تبدیل‌ها را به توابع کوچک‌تر و قابل آزمایش تقسیم کنید. علاوه بر این، از موتور قالب‌بندی Airflow (Jinja) برای انتقال مقادیر پویا بین وظایف استفاده کنید. این کار از تکرار کد می‌کاهد و نگهداری DAGها را آسان‌تر می‌سازد. برای مثال، می‌توانید مسیرهای فایل یا کوئری‌های SQL را بر اساس تاریخ اجرای فعلی به صورت پویا تولید کنید.

استقرار در محیط تولید

اجرای Airflow در محیط تولید نیازمند برنامه‌ریزی معماری دقیق است. اگرچه پایگاه داده SQLite پیش‌فرض برای توسعه مناسب است، محیط‌های تولید باید از یک پایگاه داده رابطه‌ای مقاوم مانند PostgreSQL یا MySQL برای ذخیره‌سازی متادیتا استفاده کنند. دسترس‌پذیری بالا با اجرای چندین زمان‌بند (Scheduler) و کارگر (Worker) در Airflow حاصل می‌شود. برای مقیاس‌پذیری، استفاده از KubernetesExecutor را در نظر بگیرید که به صورت پویا پاد‌های کوبرنیتس را برای هر وظیفه راه‌اندازی می‌کند و به شما اجازه می‌دهد از مدیریت منابع بومی ابری بهره ببرید. همیشه اطمینان حاصل کنید که تصاویر Docker شما سبک و قابل تکرار هستند و از پایپ‌لاین‌های CI/CD برای تست منطق DAGها پیش از استقرار استفاده کنید.

نتیجه‌گیری

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