Тяжёлый запрос на верхнем уровне файла 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 держит курсор открытым и подтягивает данные порциями, поэтому база повторяет один и тот же запрос; достаточно перейти на драйвер напрямую.
