با انفجار حجم دادهها و توزیعیابی فزایندهی زیرساختها، نیاز به هماهنگسازی قابل اعتماد، قابل مشاهده و قابل توسعهی جریانهای کاری هرگز اینقدر بزرگ نبوده است. 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 و تقویت آن برای محیط تولید، سودهای قابل توجهی در زمینه قابلیت اطمینان و کارایی عملیاتی خواهد داشت.