Ответ
Apache Airflow — это платформа с открытым исходным кодом для программируемого оркестрирования, планирования и мониторинга рабочих процессов (workflows). Её основная абстракция — DAG (Directed Acyclic Graph), который представляет собой набор задач с зависимостями.
Ключевые концепции:
- DAG: Определяет workflow как код на Python. Каждый DAG имеет расписание (cron-выражение) и дату начала.
- Operators (Операторы): Шаблоны для выполнения задач (например,
PythonOperator,BashOperator,DockerOperator). Они определяют что делать. - Tasks: Конкретные экземпляры операторов в DAG.
- Task Instances: Конкретное выполнение задачи в определённый момент времени.
Пример простого DAG для ETL-пайплайна:
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
def extract():
# Код для извлечения данных
return data
def transform(raw_data):
# Код для трансформации
return cleaned_data
def load(data_to_load):
# Код для загрузки в хранилище
pass
with DAG(
dag_id='simple_etl',
start_date=datetime(2023, 1, 1),
schedule_interval='@daily'
) as dag:
extract_task = PythonOperator(
task_id='extract',
python_callable=extract
)
transform_task = PythonOperator(
task_id='transform',
python_callable=transform,
op_args=[extract_task.output] # Зависимость по данным
)
load_task = PythonOperator(
task_id='load',
python_callable=load,
op_args=[transform_task.output]
)
extract_task >> transform_task >> load_task # Определение порядка
Типичные сценарии использования: Оркестрация ETL/ELT процессов, запуск ML-пайплайнов переобучения, автоматизация отчётов, управление инфраструктурой.