Привет, Хабр!
Меня зовут Андрей Бирюков. Я — независимый эксперт в области ИТ и ИБ, преподаю в учебных центрах и пишу статьи и книги.
С проблемой получения качественных данных сталкивался каждый аналитик, и всем понятно, что низкое качество входных данных может привести к весьма печальным результатам. Так американский ритейлер Target потерял $5,4 млрд из-за проблем с инвентаризацией, а финансовый гигант J.P. Morgan — $6,2 млрд из-за неточных данных для риск-моделей. Эти примеры ясно говорят о том, что качество данных не является абстрактной «чистотой», а представляет собой измеримую характеристику пригодности информации для конкретных бизнес-целей. Принцип GIGO («мусор на входе — мусор на выходе») в системах аналитики и ML стоит компаниям миллионов.
Как системные аналитики мы должны рассматривать качество данных не как задачу отдела СХД, а как архитектурное требование, которое закладывается на этапе проектирования. И в этой статье мы поговорим о том, как можно обеспечить получение качественных данных с практическими примерами реализации.
Шесть измерений качества
Прежде чем внедрять методы, нужно четко определить, что именно мы измеряем. Отраслевые стандарты (DAMA-DMBOK, ISO 8000) выделяют шесть измерений.

1. Полнота (Completeness). Нам важно наличие в собранных данных всех необходимых атрибутов и в качестве метрики мы определяем процент непустых значений в поле. Например, в заказе не указан ИНН клиента или не заполнено поле с телефоном.
2. Уникальность (Uniqueness). Здесь нам нужно, чтобы в собранных данных отсутствовали дубликаты записей. Метрика: процент уникальных значений относительно общего числа строк. Типовой пример это один и тот же клиент заведенный трижды с разными ID.
3. Согласованность (Consistency) представляет собой соответствие форматов и паттернов. Метрики консистентности представляют собой процент записей, соответствующих regex выражениям (email, телефон). Здесь типичный пример, когда в CRM клиент — «Физлицо», а в биллинге — «Юрлицо».
4. Валидность (Validity) — соответствие справочникам и диапазонам. Метрикой здесь является процент значений, найденных в референсном справочнике. А в качестве примера можно рассмотреть email без «@», возраст -150 лет.
5. Точность (Accuracy). А этот параметр отвечает за соответствие реальности и его измерением является процент расхождений между источником и целевой системой. Типичный пример: адрес доставки в базе не совпадает с фактическим.
6. Своевременность (Timeliness). А это актуальность данных на момент принятия решения. Здесь мы замеряем задержку обновления данных (latency). Например, отчет по продажам за вчера, сформированный сегодня в 18:00, бесполезен для утренней планёрки.
Внедрение этих метрик должно начинаться с выполнения профилирования данных. Этот процесс представляет собой автоматический анализ, который собирает статистику (распределения, мин/макс, пустые значения) и выявляет аномалии еще до написания правил очистки.
А теперь, давайте посмотрим немного кода. Для проверки полноты и уникальности на уровне колонок в Apache Airflow есть оператор SQLColumnCheckOperator, который позволяет описать проверки для каждой колонки декларативно.
Вот пример из реального пайплайна обработки заказов:
column_checks = SQLColumnCheckOperator( task_id='column_checks', table='orders', column_mapping={ 'order_id': { "unique_check": {"equal_to": 0} # не должно быть дублей }, 'price': { "min": {"greater_than": 0} # цена не может быть отрицательной }, 'quantity': { "max": {"less_than": 1000} # верхняя граница } }, conn_id=MY_CONN_ID, )
Здесь unique_check проверяет уникальность order_id, min и max валидируют диапазоны цен и количества. Если хотя бы одно условие не выполняется, оператор падает и останавливает пайплайн.
Для правил, которые связывают несколько колонок, используется SQLTableCheckOperator. В качестве при мера рассмотрим проверку финансовых показателей:
SQLTableCheckOperator( task_id="check_financial_rules", table="bookings", checks={ "total_net_fare_usd >= 0": {"partition_clause": "..."}, "total_discounts_usd <= total_gross_fare_usd": {} }, conn_id=_DB_CONN_ID, )
Этот код проверяет, является ли чистая выручка неотрицательной, а скидки не превышают валовую стоимость. Такие правила невозможно выразить на уровне отдельной колонки — нужна агрегация по строкам.
Выявляем аномалии
Теперь давайте попробуем обнаружить аномалии в полученных данных. ОператорSQLIntervalCheckOperator сравнивает метрику за сегодня с тем же показателем N дней назад, вычисляя отношение (ratio) и проверяя, не превышает ли оно заданный порог. Это действие позволяет отловить резкие скачки, которые могут указывать на проблемы с данными или реальные аномалии.
SQLIntervalCheckOperator( task_id="check_report_revenue_vs_last_week", table="daily_planet_report", date_filter_column="report_date", days_back=-7, ratio_formula="max_over_min", metrics_thresholds={"SUM(total_net_fare_usd)": 3}, ignore_zero=False, )
Здесь проверяется, что суммарная выручка не отличается от прошлой недели более чем в 3 раза. Параметр ignore_zero=False критически важен: если исторических данных нет, проверка должна упасть, а не молча пройти.

В этом примере выручка превысила 3 и в итоге наша задача завершится с ошибкой.
Проверка целостности через кастомный SQL
Еще один полезный оператор SQLCheckOperator позволяет написать произвольный запрос, который должен вернуть 1 (успех) или 0 (провал). Пример проверки наличия данных за нужный период и отсутствия NULL в критичном поле:
check_data_exist = SQLCheckOperator( task_id='check_data_exist', sql="SELECT COUNT(*) FROM orders WHERE order_date > '2023-01-01'", conn_id=MY_CONN_ID, )
Если запрос вернёт 0, оператор интерпретирует это как False и упадёт. Можно комбинировать несколько проверок через UNION ALL в одном запросе, чтобы валидировать сразу множество условий за один проход.
Обработка пропусков в данных
На этапе очистки (например, в Python/pandas) пропуски можно заполнять разными способами. Простое заполнение средним или модой (значением, которое встречается в наборе данных чаще всего) тоже работает, но может искажать распределение. Более продвинутым подходом является использование KNN-импутации, которая ищет ближайших «соседей» записи и восстанавливает пропущенное значение на основе этих значений. Например, если у пациента не указан уровень глюкозы, но известны возраст, ИМТ и давление, KNN найдёт похожих пациентов с известной глюкозой и оценит пропуск. Во многих случаях такой подход поможет исправить пропуски в данных.
От грязных данных к надёжному пайплайну
В качестве практического примера давайте построим пайплайн загрузки данных о заказах. Здесь в качестве источника будет выступать API, который может возвращать некорректные email (10% случаев), возраст с выбросами (5% — значения 0 или 120), и дубликаты записей (20 из 180 строк). В таком случае, основными шагами пайплайн будут следующие:
Шаг 1. Ingestion: сырые данные попадают в staging-таблицу без проверок.
Шаг 2. Column checks: SQLColumnCheckOperator проверяет age на диапазон (18–70), email на соответствие паттерну, order_id на уникальность.
Шаг 3. Table checks: SQLTableCheckOperator проверяет бизнес-правила (например, signup_date >= '2020-01-01').
Шаг 4. Interval checks: SQLIntervalCheckOperator сравнивает количество заказов за сегодня с прошлой неделей.
Шаг 5. Получаем результат: если критичные проверки падают, то мы останавливаем пайплайн и данные не публикуются. Если падают не критичные, то генерируется алерт, но пайплайн продолжает работу.
Пример реализации данного пайплайна на Airflow представлен ниже:
from datetime import datetime, timedelta import pandas as pd from airflow import DAG from airflow.providers.postgres.hooks.postgres import PostgresHook from airflow.providers.common.sql.operators.sql import ( SQLColumnCheckOperator, SQLTableCheckOperator, SQLIntervalCheckOperator, ) from airflow.operators.python import PythonOperator from airflow.operators.empty import EmptyOperator MY_CONN_ID = "postgres_orders" MY_DB_CONN_ID = "postgres_orders" default_args = { "owner": "data_team", "retries": 2, "retry_delay": timedelta(minutes=5), "email_on_failure": True, "email": ["data-alerts@example.com"], } with DAG( dag_id="orders_quality_pipeline", default_args=default_args, start_date=datetime(2024, 1, 1), schedule_interval="@daily", catchup=False, tags=["quality", "orders"], ) as dag: start = EmptyOperator(task_id="start") # --------------------------------------------------------------- # 1. Извлечение сырых данных из API в staging # --------------------------------------------------------------- def extract_orders(**context): hook = PostgresHook(postgres_conn_id=MY_CONN_ID) # В реальном проекте здесь был бы вызов API. # Для примера — читаем из источника-заглушки. raw_df = pd.read_sql( "SELECT * FROM source_orders WHERE order_date = %(ds)s", hook.get_conn(), params={"ds": context["ds"]}, ) raw_df.to_sql( "raw_orders", hook.get_sqlalchemy_engine(), if_exists="replace", index=False, ) extract = PythonOperator( task_id="extract_orders", python_callable=extract_orders, ) # --------------------------------------------------------------- # 2. Очистка данных: дедупликация, валидация email, отсечение выбросов # --------------------------------------------------------------- def clean_orders(**context): hook = PostgresHook(postgres_conn_id=MY_CONN_ID) engine = hook.get_sqlalchemy_engine() df = pd.read_sql("SELECT * FROM raw_orders", engine) # 2.1. Удаление дубликатов по order_id (оставляем последнюю запись) df = df.sort_values("updated_at").drop_duplicates( subset=["order_id"], keep="last" ) # 2.2. Валидация email: отбрасываем строки с некорректным форматом email_pattern = r"^[\w\.-]+@[\w\.-]+\.\w+$" valid_email_mask = df["email"].str.match(email_pattern, na=False) df = df[valid_email_mask] # 2.3. Отсечение выбросов по возрасту: оставляем 18–70 df = df[(df["age"] >= 18) & (df["age"] <= 70)] # 2.4. Заполнение пропусков в quantity медианой df["quantity"] = df["quantity"].fillna(df["quantity"].median()) df.to_sql( "clean_orders", engine, if_exists="replace", index=False, ) clean = PythonOperator( task_id="clean_orders", python_callable=clean_orders, ) # --------------------------------------------------------------- # 3. Проверка колонок (полнота, уникальность, диапазоны) # --------------------------------------------------------------- column_checks = SQLColumnCheckOperator( task_id="column_checks", table="clean_orders", column_mapping={ "order_id": { "unique_check": {"equal_to": 0}, # дублей быть не должно "null_check": {"equal_to": 0}, # NULL недопустим }, "email": { "null_check": {"equal_to": 0}, }, "age": { "min": {"greater_than_or_equal_to": 18}, "max": {"less_than_or_equal_to": 70}, }, "quantity": { "min": {"greater_than": 0}, "max": {"less_than": 1000}, }, }, conn_id=MY_DB_CONN_ID, ) # --------------------------------------------------------------- # 4. Проверка бизнес-правил на уровне таблицы # --------------------------------------------------------------- table_checks = SQLTableCheckOperator( task_id="table_checks", table="clean_orders", checks={ "net_fare_not_negative": { "check_statement": "total_net_fare_usd >= 0" }, "discounts_not_exceed_gross": { "check_statement": "total_discounts_usd <= total_gross_fare_usd" }, "row_count_positive": { "check_statement": "COUNT(*) > 0" }, }, conn_id=MY_DB_CONN_ID, ) # --------------------------------------------------------------- # 5. Сравнение с историческими данными (аномалии) # --------------------------------------------------------------- interval_checks = SQLIntervalCheckOperator( task_id="check_report_revenue_vs_last_week", table="daily_planet_report", date_filter_column="report_date", days_back=-7, ratio_formula="max_over_min", metrics_thresholds={"SUM(total_net_fare_usd)": 3}, ignore_zero=False, ) # --------------------------------------------------------------- # 6. Публикация в целевую витрину # --------------------------------------------------------------- def publish_orders(**context): hook = PostgresHook(postgres_conn_id=MY_CONN_ID) engine = hook.get_sqlalchemy_engine() df = pd.read_sql("SELECT * FROM clean_orders", engine) df.to_sql( "published_orders", engine, if_exists="append", index=False, ) publish = PythonOperator( task_id="publish_orders", python_callable=publish_orders, ) end = EmptyOperator(task_id="end") # --------------------------------------------------------------- # Оркестрация: жёсткие шлюзы для критичных проверок # --------------------------------------------------------------- start >> extract >> clean >> column_checks >> table_checks >> interval_checks >> publish >> end
Airflow позволяет по разному реагировать на провалы проверок при обработке ошибок качества данных. При использовании варианта Жёсткий шлюз (Hard gate), проверки встроены в основной DAG и если они завершаются с ошибкой, то весь DAG также падает, и downstream-потребители не получают данные.
Вариант Средний шлюз (Medium gate) предполагает использование отдельного DAG для проверок, который запускается после обработки. Если проверки падают, DAG тоже падает, но при этом, основной пайплайн уже завершился. В результате проблема становится видимой, но данные уже доступны.
И вариант с Мягким шлюзом (Soft gate) предполагает, что проверки встроены, но не влияют на финальный статус DAG. Реализуется через trigger_rule="all_done" на последней задаче, то есть алерты срабатывают, но пайплайн считается успешным. Здесь важно понимать, что это достаточно опасный подход, потому что реальные ошибки других задач тоже могут быть проигнорированы.
Подведем итоги
Обеспечение качества данных представляет собой комбинацию профилирования, декларативных правил, устранения противоречий и наблюдаемости. Современные инструменты вроде Airflow SQL Check Operators позволяют встроить проверки прямо в пайплайн с минимальной нагрузкой. И здесь главное правило заключается в том, что проверки должны быть жёсткими для критичных данных и предупреждающими — для второстепенных.

Качество данных влияет на решения, которые принимаются на их основе: от аналитических отчётов до работы ML-систем. На бесплатных уроках можно подробнее разобрать подходы к управлению качеством данных, познакомиться с практическими кейсами экспертов и обсудить вопросы, которые возникают при построении таких процессов в командах.
30 сентября в 20:00. «Качество данных: риски и ответственность». Записаться
13 октября в 20:00. «Технологии и принципы построения DQaaS (Data Quality as a Service) в современных системах управления данными». Записаться
20 октября в 20:00. «Метрики качества данных и стратегия внедрения». Записаться