Как связать процессы в пайплайн?
В первой статье показал общую схему проекта eltdwh-airflow-dbt. Как мы уже выяснили, главный принцип устойчивой системы — не превращать пайплайн в монолитный скрипт, а дробить процесс на короткие автономные операции.
Но когда у вас появляется десяток таких операций, встает вопрос: кто будет ими управлять?
Кто проверит, что первый шаг завершился успешно, прежде чем запускать второй?
Кто сделает выбор, какой процесс за каким должен следовать?
Для этого существует оркестратор (в проекте использую Apache Airflow). Его задача — не выполнять саму работу, а быть умным диспетчером.
Давайте посмотрим на логику первого DAG-а проекта (1_precheck_xlsx_intake). Он выполняет роль «швейцара» на входе в систему и состоит из 5 коротких тасок:
1. Сканируем папку на наличие файлов: t1_scan_dir >> t2_branch_if_exists 2. Если файлов нет — мгновенно завершаем процесс (не тратим ресурсы): t2_branch_if_exists >> finish_no_files 3. Если файлы есть — запускаем цепочку проверок: t2_branch_if_exists >> t3_classify_files >> t4_check_idempotency >> t5_decide_trigger 4. Если файлы уже загружались раньше — выходим: t5_decide_trigger >> finish_no_files 5. И только если всё корректно — даем команду на загрузку следующему DAG: t5_decide_trigger >> trigger_ingest_xlsx
В чем здесь чисто инженерная красота и польза?
• Умное ветвление (BranchPythonOperator): Пайплайн запускается каждую минуту. Благодаря таске t2, если новых данных нет, оркестратор завершает выполнение. Мы не тратим ресурсы зря. • Изоляция ошибок: таска t3 (классификация) разделяет данные. Если в папке лежит битый файл, она переместит его в директорию error, но не даст упасть всему пайплайну — корректные данные пойдут дальше. • Идемпотентность (таска t4): Перед тем как пустить данные в базу, мы сверяем имя файла с данными в PostgreSQL. Если файл уже обрабатывался — он отбрасывается. Это гарантия того, что данные в вашем DWH не задвоятся при повторном запуске.
Главный вывод:
Каждая из этих тасок внутри Airflow — это буквально 15–30 строк понятного Python-кода. Их легко читать и дебажить. Но за счет связей между ними мы получаем гибкую, отказоустойчивую систему, которая умеет принимать решения на ходу.
Оркестратор выстроил цепочку, проверил безопасность и доставил файл к воротам базы данных. Теперь начинается трансформация. И здесь важно не совершить главную ошибку — не превратить Airflow в «чернорабочего», который сам крутит тяжелый SQL.
О том, почему оркестратор должен только отдавать приказы, а всю тяжелую работу в базе должен делать dbt, поговорим в следующей статье.
Ссылка на код DAG-а — в первом комментарии.
#DataEngineering #Airflow #BestPractices #Architecture #Backend #eltdwh_pipeline #инженерияданных
· 12.08
Код этого DAG-а и функции валидации выложены в репозитории: https://github.com/lelik-bolek/eltdwh-airflow-dbt
0
ответить
коммент скрыт — часть юзеров считает его токсичным или некорректным
коммент удалён