DAG в Airflow — это обычный Python-файл, который описывает задачи и порядок их запуска. В этом уроке напишем DAG из двух шагов, положим его в нужную папку и посмотрим, как планировщик подхватит его без перезапуска.
Структура файла
Три обязательных части: объект DAG, хотя бы один оператор и зависимости между ними.
from airflow import DAG
from airflow.operators.bash import BashOperator
with DAG("hello_de", schedule="@daily") as dag:
extract = BashOperator(task_id="extract", bash_command="echo extract")
load = BashOperator(task_id="load", bash_command="echo load")
extract >> loadОператор >> — это перегруженный сдвиг: читается как «сначала extract, потом load».
Кладём DAG в папку
Airflow сканирует папку dags/ каждые несколько секунд. Достаточно сохранить файл, и через минуту DAG появится в интерфейсе.
Запуск и логи
Нажмите на DAG, включите тумблер слева от названия и запустите вручную кнопкой Trigger. Каждая задача пишет свой лог, его можно открыть кликом по квадратику.
Типичные ошибки
- Файл лежит не в той папке.
- В
task_idесть пробелы или кириллица. - Два DAG с одинаковым
dag_id.