Вступить в клуб →
сложнаявопросОркестрация пайплайнов в Airflow

Тяжёлый запрос на верхнем уровне файла DAG нагружает источник

После выкатки этого DAG планировщик стал заметно тормозить, а DBA прислал график: к базе-источнику каждые несколько секунд приходит одинаковый тяжёлый SELECT — хотя DAG запускается раз в сутки.

# dags/clients_dag.py
import pandas as pd

df = pd.read_sql("select distinct client_id from dds.clients", conn)
client_ids = df.client_id.tolist()

with DAG("clients_dag", schedule="@daily", ...) as dag:
    for cid in client_ids:
        PythonOperator(task_id=f"load_{cid}", python_callable=load, op_args=[cid])

В чём причина и что делать?

  • Запрос выполняется заново при старте каждой задачи: воркер импортирует файл DAG перед исполнением. Нужно вынести чтение в отдельную задачу и передать список дальше через XCom целиком.
  • Планировщик регулярно перечитывает файлы DAG-ов, поэтому код верхнего уровня исполняется при каждом парсинге; тяжёлое чтение нужно унести внутрь задачи, а список сущностей брать из дешёвого источника.
  • Проблема в динамическом создании задач как таковом: число задач в DAG обязано быть фиксированным, иначе планировщик считает структуру изменившейся и пересоздаёт DAG на каждом тике.
  • Причина в pandas: read_sql держит курсор открытым и подтягивает данные порциями, поэтому база повторяет один и тот же запрос; достаточно перейти на драйвер напрямую.

🔒 Проверка ответа — для участников клуба

  • Проверка ответа
  • Подсказка, если застряли
  • Разбор с объяснением, почему так
  • Прогресс по всем задачам и виртуальные собеседования
Зарегистрироваться →

Регистрация занимает минуту

Другие задачи раздела