Ответ
Да, я реализовывал асинхронные DQ-проверки для пайплайнов, где немедленная остановка потока была нежелательна.
Пример из проекта: У нас был Airflow DAG, который загружал ежедневные пакеты транзакций. Вместо встроенных проверок в основном DAG, мы создали отдельный DAG для качества данных.
- Основной DAG (
load_transactions) загружал данные в staging-таблицу и по завершению публиковал событие (например, записывал метаданные в служебную таблицу или отправлял сообщение в Kafka). - DQ DAG (
check_transactions_quality) запускался по триггеру от этого события и выполнял проверки:# Пример проверки в Python (Great Expectations или кастомный код) def check_row_count(): # Сравнение количества загруженных строк с ожидаемым pass def check_duplicates(): # Поиск дубликатов по ключевым полям pass def check_value_ranges(): # Проверка, что суммы транзакций положительные pass - Результаты записывались в отдельную таблицу
dq_resultsи отправлялись алертами в Slack/Email, если были найдены критические нарушения.
Преимущество: Основной процесс загрузки не блокировался, и мы могли гибко настраивать пороги срабатывания алертов и глубину проверок.