Настройка маршрутизации процессов в Apache Airflow: пошаговое руководство | AdminWiki

Настройка маршрутизации процессов в Apache Airflow: пошаговое руководство

22 июля 2026 10 мин. чтения
Содержание статьи

Маршрутизация в Apache Airflow - это логика, которая определяет порядок выполнения задач внутри DAG (Directed Acyclic Graph). Вы задаете правила: какие задачи запускать последовательно, какие параллельно, а какие - только при соблюдении внешних условий. Этот механизм опирается на операторы (действия) и сенсоры (ожидания), а также на явные зависимости между задачами, прописанные в коде. Правильно спроектированная маршрутизация исключает гонки данных, позволяет обрабатывать сбои без ручного вмешательства и делает пайплайн прозрачным для мониторинга.

В этом руководстве разобраны ключевые инструменты управления потоком: от базовых цепочек task1 >> task2 до условного ветвления через BranchPythonOperator и синхронизации между DAG с помощью ExternalTaskSensor. Вы получите готовые шаблоны для ETL-пайплайнов, настройки повторных попыток и интеграции с системами мониторинга. Все примеры проверены на актуальных версиях Apache Airflow 2026 года.

Если вы проектируете новый пайплайн или оптимизируете существующий, понимание принципов маршрутизации сократит время отладки и повысит надежность обработки данных. Детали архитектуры современных оркестраторов и сравнение Airflow с Dagster и Prefect рассмотрены в отдельном материале по ETL/ELT-пайплайнам.

Основы маршрутизации: как Airflow управляет потоком задач

DAG в Airflow - это направленный ациклический граф, где узлы - задачи, а ребра - зависимости. Планировщик анализирует граф и запускает задачу только тогда, когда все её родительские задачи успешно завершены. Этот принцип лежит в основе всей маршрутизации.

Зависимости задаются битовыми операторами >> и <<. Airflow интерпретирует их как направленные связи. Когда вы пишете task1 >> task2, вы указываете: task2 не начнется, пока task1 не выполнится. Для параллельного запуска используется список: task1 >> [task2, task3] запустит task2 и task3 одновременно после завершения task1.

Планировщик проверяет состояние всех upstream-задач. Если хотя бы одна из них находится в статусе failed, downstream-задача по умолчанию получает статус upstream_failed и не запускается. Это поведение можно изменить через trigger_rule, но базовая логика гарантирует, что данные не будут обработаны до готовности источника.

Зависимости задач: последовательное и параллельное выполнение

Линейная цепочка extract >> transform >> load - простейший маршрут. Каждый этап ждет завершения предыдущего. Для ускорения пайплайна трансформацию можно распараллелить:

extract >> [transform_users, transform_orders] >> load

Здесь обе задачи трансформации запускаются одновременно. Задача load начнется только после успешного завершения обеих. Airflow определяет готовность задачи к запуску, проверяя статусы всех прямых предков. Если transform_users завершился успешно, а transform_orders еще выполняется, load будет ждать.

Для объединения результатов параллельных веток используйте пустую задачу-сборщик с trigger_rule='ALL_SUCCESS'. Это явно указывает планировщику дождаться всех веток, даже если одна из них пропущена по условию.

Операторы и сенсоры: два столпа управления потоком

Операторы выполняют работу. PythonOperator запускает произвольный код, BashOperator - команды оболочки, PostgresOperator - SQL-запросы. Каждый оператор инкапсулирует одно атомарное действие. Выбор оператора определяет, что именно произойдет на этом шаге маршрута.

Сенсоры ждут внешние условия. FileSensor проверяет появление файла в файловой системе, ExternalTaskSensor - завершение задачи в другом DAG, TimeDeltaSensor - истечение заданного интервала. Сенсор занимает слот работника, пока условие не выполнится или не истечет таймаут. Для длительных ожиданий используйте режим mode='reschedule', чтобы освободить слот между проверками.

Комбинация операторов и сенсоров позволяет строить маршруты, которые адаптируются к внешней среде: ждут поступления данных, реагируют на завершение смежных процессов и выполняют полезную нагрузку.

Ветвление и условное выполнение задач

Линейные маршруты не покрывают сценарии, где путь выполнения зависит от данных. Airflow предоставляет два инструмента для условного ветвления: BranchPythonOperator для выбора одной ветки из нескольких и ShortCircuitOperator для досрочного завершения цепочки.

Оба оператора принимают Python-функцию, которая возвращает идентификатор следующей задачи или булево значение. Разница в поведении: BranchPythonOperator пропускает невыбранные ветки, ShortCircuitOperator останавливает всю downstream-цепочку.

BranchPythonOperator: принятие решений в пайплайне

Функция для BranchPythonOperator должна вернуть task_id следующей задачи или список идентификаторов для нескольких веток. Задачи, которые не попали в возвращаемый список, получают статус skipped.

from airflow.operators.python import BranchPythonOperator

def choose_path(**kwargs):
    ti = kwargs['ti']
    data_quality = ti.xcom_pull(task_ids='check_quality')
    if data_quality > 0.95:
        return 'load_to_production'
    else:
        return 'load_to_staging'

branch = BranchPythonOperator(
    task_id='branch',
    python_callable=choose_path
)

Задача, следующая за точкой ветвления, должна иметь trigger_rule='ALL_DONE' или 'ONE_SUCCESS'. Правило ALL_SUCCESS (по умолчанию) не сработает, потому что одна из веток будет пропущена. ALL_DONE запустит задачу, когда все upstream-задачи перейдут в финальное состояние - success, skipped или failed.

ShortCircuitOperator: прерывание ветки по условию

ShortCircuitOperator вычисляет условие и, если оно ложно, пропускает все downstream-задачи в своей ветке. Это удобно для проверки наличия данных перед ресурсоемкой обработкой.

from airflow.operators.python import ShortCircuitOperator

def has_new_data(**kwargs):
    # Проверяем наличие файлов в источнике
    files = list_new_files('/data/incoming')
    return len(files) > 0

check_data = ShortCircuitOperator(
    task_id='check_data',
    python_callable=has_new_data
)

check_data >> transform >> load

Если has_new_data возвращает False, задачи transform и load получают статус skipped. В отличие от BranchPythonOperator, здесь не нужно настраивать trigger_rule для downstream-задач - они просто не запускаются.

Ожидание внешних событий с помощью сенсоров

Сенсоры синхронизируют DAG с внешним миром. Они периодически проверяют условие с интервалом poke_interval (по умолчанию 60 секунд) и удерживают слот работника до наступления события или истечения timeout. Для длительных ожиданий переключайте сенсор в режим mode='reschedule': между проверками слот освобождается, и работник может выполнять другие задачи.

ExternalTaskSensor: зависимость между DAG

Этот сенсор ждет завершения задачи в другом DAG. Он критичен для построения цепочек пайплайнов, где один процесс запускается строго после другого.

from airflow.sensors.external_task import ExternalTaskSensor

wait_for_ingestion = ExternalTaskSensor(
    task_id='wait_for_ingestion',
    external_dag_id='data_ingestion',
    external_task_id='load_raw_data',
    timeout=3600,
    poke_interval=120,
    mode='reschedule'
)

Основная сложность - синхронизация по execution_date. По умолчанию сенсор ищет запуск внешнего DAG с той же датой выполнения. Если расписания DAG не совпадают, используйте параметры execution_delta или execution_date_fn для сдвига. Параметр window_size позволяет искать в диапазоне дат, что полезно при рассинхронизации расписаний.

TimeDeltaSensor и расписания: временная маршрутизация

TimeDeltaSensor вставляет фиксированную паузу в пайплайн. Это нужно, например, чтобы дождаться завершения фоновых процессов во внешней системе, у которых нет API для проверки статуса.

from airflow.sensors.time_delta import TimeDeltaSensor

wait_30_minutes = TimeDeltaSensor(
    task_id='wait_30_minutes',
    delta=timedelta(minutes=30)
)

Расписание DAG задается параметром schedule_interval. Оно определяет, когда Airflow создает новый экземпляр DAG (DAG Run). Для ежедневного запуска в 3:00 утра используйте cron-выражение '0 3 * * *'. Расписание влияет на маршрутизацию косвенно: каждый DAG Run получает свою execution_date, которая используется сенсорами для синхронизации.

Обработка ошибок и повышение отказоустойчивости

Сбои в распределенных системах неизбежны. Airflow компенсирует их автоматическими повторными попытками, правилами триггеров и колбэками для оповещения. Настройка этих механизмов превращает хрупкий скрипт в промышленный пайплайн, способный восстанавливаться без участия дежурного инженера.

Повторные попытки и задержки: настройка retries

Параметры retries и retry_delay задают количество повторных запусков и паузу между ними. Настройка применяется на уровне задачи или через default_args для всего DAG.

default_args = {
    'retries': 3,
    'retry_delay': timedelta(minutes=5),
    'retry_exponential_backoff': True,
    'max_retry_delay': timedelta(minutes=30)
}

Экспоненциальная задержка (retry_exponential_backoff=True) увеличивает паузу после каждой неудачной попытки: 5, 10, 20 минут. Это снижает нагрузку на внешнюю систему, если сбой вызван её перегрузкой. Параметр max_retry_delay ограничивает максимальную паузу.

Правила триггеров: управление потоком после сбоя

trigger_rule определяет, при каких состояниях upstream-задач запускается текущая. По умолчанию действует ALL_SUCCESS - задача ждет успеха всех предков. Для обработки ошибок важны следующие правила:

  • ALL_FAILED - запуск, если все upstream-задачи упали. Используется для задач очистки после полного сбоя.
  • ALL_DONE - запуск в любом случае, независимо от статуса предков. Подходит для отправки уведомлений или освобождения ресурсов.
  • ONE_SUCCESS - запуск, если хотя бы одна задача выполнилась успешно. Полезно при параллельных запросах к нескольким репликам данных.
  • ONE_FAILED - запуск при первом же сбое. Позволяет быстро отреагировать на проблему, не дожидаясь остальных задач.
cleanup = BashOperator(
    task_id='cleanup',
    bash_command='rm -rf /tmp/processing/*',
    trigger_rule='ALL_DONE'
)

Задача cleanup выполнится всегда - и после успеха, и после падения основного пайплайна. Это гарантирует, что временные файлы не забьют диск.

Колбэки on_failure_callback и on_retry_callback вызываются при соответствующих событиях. Типичный сценарий - отправка уведомления в Slack или email при падении задачи. Функция колбэка получает контекст с информацией о задаче, логах и причине сбоя.

Оптимизация производительности и масштабирование

Производительность пайплайна определяется тем, насколько эффективно Airflow распределяет задачи по работникам и управляет доступом к внешним ресурсам. Три ключевых механизма: параллелизм, пулы и приоритеты задач.

Управление параллелизмом и пулы ресурсов

Параметр parallelism в конфигурации Airflow ограничивает общее количество одновременно выполняющихся задач во всех DAG. Значение по умолчанию - 32. Для высоконагруженных инсталляций его увеличивают до 100-200, но с оглядкой на ресурсы базы данных и количество работников.

Пулы (pools) ограничивают параллелизм на уровне конкретного ресурса. Например, если внешняя база данных выдерживает только 5 одновременных подключений, создайте пул с slots=5 и назначьте его всем задачам, работающим с этой базой.

transform = PythonOperator(
    task_id='transform',
    python_callable=transform_data,
    pool='postgres_pool',
    priority_weight=10
)

Параметр priority_weight управляет очередью выполнения. Задачи с большим весом получают слоты раньше. Это полезно, когда критичные пайплайны не должны ждать фоновых задач. Веса можно назначать динамически через priority_weight_fn.

Динамические DAG: генерация задач во время выполнения

Когда количество задач зависит от входных данных, DAG генерируется динамически. Типичный пример - обработка переменного числа файлов: для каждого файла создается отдельная задача.

def create_dag(dag_id, schedule, file_list):
    dag = DAG(dag_id, schedule_interval=schedule)
    
    with dag:
        start = DummyOperator(task_id='start')
        end = DummyOperator(task_id='end')
        
        for file_name in file_list:
            task = PythonOperator(
                task_id=f'process_{file_name}',
                python_callable=process_file,
                op_kwargs={'file': file_name}
            )
            start >> task >> end
    
    return dag

Важно: список задач фиксируется в момент парсинга DAG-файла. Изменение file_list между запусками требует повторного парсинга. Для полностью динамических сценариев, где задачи определяются во время выполнения, используйте TaskGroup с динамическим маппингом или подключайте внешние конфигурации через Variables и Connections.

Мониторинг и алертинг: контроль состояния пайплайнов

Мониторинг в Airflow строится на трех уровнях: встроенные метрики, SLA-соглашения и интеграция с внешними системами. Подробное руководство по настройке полного конвейера телеметрии с визуализацией в Grafana доступно в статье по маршрутизации логов и метрик.

SLA и уведомления о задержках

SLA (Service Level Agreement) в Airflow - это максимальное допустимое время выполнения задачи. Если задача не завершилась до истечения SLA, Airflow отправляет email-уведомление и регистрирует событие в логах.

default_args = {
    'sla': timedelta(hours=1),
    'email': ['devops-team@example.com'],
    'email_on_failure': True
}

SLA срабатывает независимо от статуса задачи. Даже если задача еще выполняется, но уже превысила лимит, уведомление будет отправлено. Это позволяет обнаружить зависшие задачи до того, как они заблокируют весь пайплайн.

Интеграция с внешними системами мониторинга

Airflow отдает метрики через statsd. Настройте statsd_on=True в airflow.cfg и укажите хост и порт statsd-сервера. Prometheus забирает метрики через statsd_exporter, Grafana визуализирует их на дашбордах. Ключевые метрики: количество запущенных задач, длительность выполнения DAG, количество сбоев, загрузка пулов.

Для отправки алертов в Slack или Telegram используйте колбэки. Функция получает контекст задачи и формирует сообщение с деталями сбоя. Готовые шаблоны алертов и дашбордов для SRE-практик рассмотрены в руководстве по наблюдаемости высоконагруженных систем.

Шаблоны надежных пайплайнов: примеры из практики

Теоретические знания собираются в рабочие конструкции. Два шаблона покрывают большинство сценариев обработки данных: классический ETL с обработкой ошибок и динамическое ветвление для переменного объема входных данных.

ETL-пайплайн с обработкой ошибок

Этот шаблон объединяет сенсоры, операторы и настройки отказоустойчивости в целостный пайплайн. Он ждет появления файла, проверяет качество данных, выбирает путь загрузки и обрабатывает сбои на каждом этапе.

from airflow import DAG
from airflow.operators.python import PythonOperator, BranchPythonOperator
from airflow.sensors.filesystem import FileSensor
from datetime import datetime, timedelta

default_args = {
    'retries': 2,
    'retry_delay': timedelta(minutes=5),
    'on_failure_callback': send_alert_to_slack
}

with DAG('etl_pipeline', default_args=default_args, schedule_interval='@daily') as dag:
    
    wait_file = FileSensor(
        task_id='wait_file',
        filepath='/data/incoming/daily_export.csv',
        timeout=600,
        mode='reschedule'
    )
    
    transform = PythonOperator(
        task_id='transform',
        python_callable=transform_data
    )
    
    def check_quality(**kwargs):
        score = kwargs['ti'].xcom_pull(task_ids='transform')
        return 'load_prod' if score > 0.95 else 'load_staging'
    
    branch = BranchPythonOperator(
        task_id='branch',
        python_callable=check_quality
    )
    
    load_prod = PythonOperator(task_id='load_prod', python_callable=load_to_production)
    load_staging = PythonOperator(task_id='load_staging', python_callable=load_to_staging)
    
    notify = PythonOperator(
        task_id='notify',
        python_callable=send_completion_notice,
        trigger_rule='ALL_DONE'
    )
    
    wait_file >> transform >> branch
    branch >> [load_prod, load_staging] >> notify

Ключевые моменты: FileSensor в режиме reschedule не блокирует работника на 10 минут ожидания. BranchPythonOperator выбирает целевую базу на основе метрики качества. Задача notify с trigger_rule='ALL_DONE' отправляет уведомление при любом исходе.

Пайплайн с динамическим ветвлением

Шаблон для сценариев, где количество обрабатываемых объектов определяется во время выполнения. Список файлов извлекается из внешнего источника, и для каждого файла создается задача обработки.

def get_file_list(**kwargs):
    # Запрос к API или сканирование директории
    files = scan_s3_bucket('raw-data')
    kwargs['ti'].xcom_push(key='file_list', value=files)
    return files

def process_single_file(file_name, **kwargs):
    # Обработка одного файла
    pass

with DAG('dynamic_processing', schedule_interval='@daily') as dag:
    
    scan = PythonOperator(
        task_id='scan_files',
        python_callable=get_file_list
    )
    
    # Динамическое создание задач через TaskGroup или цикл
    # В Airflow 2.6+ доступен dynamic task mapping
    process = PythonOperator.partial(
        task_id='process_file',
        python_callable=process_single_file
    ).expand(op_kwargs=[{'file_name': f} for f in get_file_list()])

Dynamic task mapping (доступен с Airflow 2.3+) позволяет создавать задачи на основе данных, полученных во время выполнения. Количество задач определяется результатом функции get_file_list. Это устраняет необходимость пересоздавать DAG при изменении списка файлов.

Принципы построения отказоустойчивых пайплайнов, включая защиту от OOM и настройку cgroups для задач-воркеров, детально разобраны в практическом руководстве по надежности и масштабируемости.

Поделиться:
Сохранить гайд? В закладки браузера