Ни для кого не секрет, что проверки качества – фундамент доверия к данным. Пока вы не знаете, насколько полны и корректны ваши данные, любые основанные на них выводы могут быть ошибочными. Важно не только проверить данные, но и сделать результаты этих проверок доступными всем заинтересованным пользователям.

Меня зовут Микаэл Новиков, я главный разработчик в команде ML-платформы MAGNIT TECH. Расскажу, как мы в проекте F&R не смогли оставить data quality без внимания, отобразили результаты проверок в OpenMetadata и попутно наступили на несколько грабель интеграции open source-решения.

Качество данных

Качество данных мы проверяем с помощью библиотеки Deequ – одного из распространённых инструментов для этой задачи в экосистеме Apache Spark. Deequ позволяет декларативно задать ограничения на данные (constraints) и прогнать по ним проверки. На выходе получается структурированный результат: список ограничений с их статусами. Поначалу результаты проверок просто складывались в PostgreSQL.

Затем мы внедрили каталог данных OpenMetadata (OMD), который довольно быстро прижился. Информация об источниках, таблицах, lineage и прочие метаданные собраны в одном месте. Логично было разместить рядом и информацию о качестве данных. Именно так мы и решили поступить.

Важно отметить, что сейчас у нас развернута OpenMetadata версии 1.8.7. Каталог данных активно развивается. Недавно было мажорное обновление. Учитывайте это!

Итак, проверки качества данных у нас уже были. И развернутый каталог данных тоже. А вот интеграции — нет.

Как интегрировать-то будем?

OpenMetadata позволяет запускать встроенные проверки качества данных. Тесты можно создать через интерфейс, а исполняются они непосредственно в OMD. Такой вариант нам не подходил, ведь мы уже имели результаты проверок, запускаемых с помощью Spark.

Небольшое исследование показало, что в каталоге данных нет возможности «из коробки» подключить внешние проверки. Зато есть открытый REST API. Его модель качества данных оказалась достаточно гибкой. Если коротко, она состоит из трёх основных сущностей:

  • TestDefinition — тип теста;

  • TestCase — конкретный экземпляр теста, привязанный к таблице или столбцу;

  • TestCaseResult — результат выполнения теста в конкретный момент времени.

Никто не мешает создать их через API: завести собственные TestDefinition под наши типы проверок Deequ, для каждой проверки создать TestCase и добавлять в него TestCaseResult по мере поступления новых результатов.

Так и родилась идея MetaBridge — сервиса, который читает результаты работы Deequ и раскладывает их по сущностям OMD.

Эволюция архитектуры MetaBridge

Первая попытка. PostgreSQL как источник

Изначальная реализация была, прямо скажем, примитивной. Поскольку результаты проверок хранились в PostgreSQL, мы, недолго думая, решили, что новый сервис должен просто читать результаты из базы и отправлять их в OMD.

К сожалению, у API OMD нет методов для добавления результатов батчами, поэтому всё, что мы могли сделать, это читать данные из БД и последовательно добавлять в OMD. При этом оказалось, что API каталога данных отвечает не быстро (в среднем около секунды на запрос). Один поток в таком режиме банально не справлялся с нагрузками. Мы начали изучать проблему, и с помощью простых тестов поняли, что API может достаточно шустро работать в асинхронном режиме. Получилось ускорить добавление данных по крайней мере в 10 раз. Хорошо, тогда читаем батчами по 10 результатов, обрабатываем, затем асинхронным HTTP-клиентом отправляем их в OMD.

Подход работал, но имел серьезные недостатки. Сервис должен напрямую опрашивать операционную базу команды инженеров данных. Чтобы не обрабатывать повторно одни и те же строки, нам приходилось бы самим помнить, что уже отправлено в OMD: вести курсор, помечать или удалять обработанные записи. По сути, мы добавляли бы в свой сервис всю логику отслеживания прогресса. Любая обработка сбоев и повторных попыток тоже целиком ложилась на нас. Добавим сюда зависимость от чужой схемы данных.

Оставался нерешенным и вопрос масштабирования. Что делать, когда одного процесса перестанет хватать? Несколько воркеров могли бы читать из общей таблицы PostgreSQL; технически это возможно, например через SELECT ... FOR UPDATE SKIP LOCKED. Но это превращение таблицы в очередь, со всеми сопутствующими хлопотами: блокировками, пометкой или удалением обработанных строк, повторной обработкой при сбоях, мониторингом отставания. Получается просто брокер, собранный поверх базы данных.

Конечно, мы не стали превращать PostgreSQL в брокер, а взяли проверенное решение — Kafka. Благо она уже была в инфраструктуре, поэтому не пришлось лишний раз тревожить коллег согласованиями и развертыванием.

Kafka поможет

Сразу отмечу, что для текущего потока данных одного воркера, читающего из БД батчами, вполне хватило бы. Важно, что Kafka мы выбрали не столько ради скорости — ее, как разобрались выше, дает асинхронная работа с API. Брокер решал и другие задачи.

Предлагаемая архитектура позволяет легко масштабировать систему при увеличении потока результатов проверок. Кроме того, используя общий брокер, сервис перестает зависеть от чужой базы данных. Теперь продюсер публикует результат анализа в топик, а MetaBridge становится обычным консьюмером.

Есть и взгляд в будущее. Результаты проверок качества данных можно считать самостоятельным data product. Сегодня этот топик читаем только мы, а завтра к нему сможет подключиться любой другой консьюмер, например, дашборд.

Для общего кластера Kafka есть требование использовать Protobuf ради экономии места на диске. Сообщения передаются в формате Protobuf со стандартным wire format от Confluent, а схемы хранятся в Schema Registry. Это дает нам строгий контракт между сервисами, о безусловной важности которого мы здесь говорить не будем.

Сам формат сообщения с результатом анализа выглядит так:

message DataQualityConstraint {
  string name = 1;
  string status = 2;
  string message = 3;
}

message DataQualityAnalysisResult {
  string dq_id = 1;
  string job_id = 2;
  string table = 3;
  string query = 4;
  string s3_path = 5;
  google.protobuf.Timestamp start_time = 6;
  google.protobuf.Timestamp end_time = 7;
  bool passed = 8;
  repeated DataQualityConstraint constraints = 9;
}

Один DataQualityAnalysisResult описывает результат анализа одной таблицы с произвольным количеством проверок внутри.

Топик

Партиции

Продюсеры

Консьюмеры

data_quality.analysis_results

10

Конвейеры проверок качества данных

MetaBridge

data_quality.transfer_to_open_metadata.dlq

1

MetaBridge

Скрипт обработки ошибок (MetaBridge)

data_quality.transfer_to_open_metadata.status

1

MetaBridge

Отсутствуют. Потенциально: сервисы мониторинга, фронтенд MetaBridge

В data_quality.analysis_results приходят сами результаты анализа, в data_quality.transfer_to_open_metadata.dlq MetaBridge публикует итог обработки каждого сообщения, а сообщения, которые обработать не удалось, отправляются в dead letter queue в data_quality.transfer_to_open_metadata.dlq. Читаем результаты мы с семантикой at-least-once: оффсет коммитится только после того, как сообщение полностью обработано. Это позволяет быть уверенным, что результат будет добавлен в OpenMetadata или отправлен в DLQ. С другой же стороны, при сбое или перезапуске одно и то же сообщение может прийти повторно — нужна идемпотентная обработка. Как именно мы с этим справляемся, разберём ниже.

Особенности реализации

Основным стеком нашей команды является Python. Был выбран асинхронный фреймворк для брокеров FastStream. Хотелось с одной стороны воспользоваться преимуществом использования фреймворка (автогенерируемая документация, удобное тестирование, DI, готовые интеграции и пр.), а с другой — попробовать новую для нас технологию, убедиться в ее зрелости.

Кстати, настоятельно советую присмотреть к FastStream всем, кто работает с Kafka, RabbitMQ, NATS, и Redis (особенно, если у вас в проекте есть разные типы брокеров). В этой статье мы немного поговорим  и о преимуществах использования этого фреймворка.

Схемы и обработка сообщений

Обычно для работы с Protobuf .proto-файлы лежат в репозитории, а сгенерированные из них классы хранятся рядом. У нас же единственным источник правды о схеме является Schema Registry. Такой подход избавляет от огромного количества проблем, особенно учитывая, что с Kafka работают много команд. Не нужно держать копии .proto-файлов в репозиториях и следить, чтобы они не разъехались.

MetaBridge генерирует Python-классы при запуск приложения. Происходит это в два простых шага. Сначала сервис читает нужные схемы по идентификаторам из Schema Registry и создает .proto-файлы. Затем запускает grpc_tools.protoc, который превращает их в *_pb2.py-модули.

Фреймворк FastStream позволяет, в первую очередь, упростить получение и отправку сообщений. Весь код для работы с брокером сокращается до нескольких строк, что особенно ценно, когда обработчиков и топиков становится много.

Вот как выглядит процесс обработки каждого сообщения:

@router.subscriber(
   config.ANALYSIS_RESULTS_TOPIC,
   group_id=config.ANALYSIS_RESULTS_CONSUMER_GROUP_ID,
   auto_offset_reset='earliest',
   ack_policy=AckPolicy.MANUAL,
   # Стандартного интервала в 5 минут иногда не хватает для обработки одного сообщения.
   max_poll_interval_ms=10 * 60 * 1000,
)
async def handle_data_quality_analysis_result(message: KafkaMessage, logger: Logger) -> None:
   # Десериализация сообщения.
   _, _, payload = parse_confluent_protobuf_message(message.body)
   result = DataQualityAnalysisResult()
   try:
       result.ParseFromString(payload)
   except DecodeError as e:
       raise ProtobufMessageDeserializationError(e) from e


   # Обработка результата проверки и добавление в OpenMetadata.
   logger.info('Adding result with dq_id=%s to Open Metadata...', result.dq_id)
   added, details = await process_data_quality_analysis_result(result)
   logger.info('Process complete: added=%s, details=%s', added, details)


   # Запись статуса добавления и, если нужно, отправка в DLQ.
   time = datetime.now(UTC)
   message_key = result.dq_id.encode('utf-8')
   if added:
       status = DataQualityTransferToOpenMetadataStatus(
           dq_id=result.dq_id,
           time=time,
           success=True,
           details=details,
       )
       await status_publisher.publish(
           message=generate_message_header(config.STATUS_SCHEMA_ID) + status.SerializeToString(),
           key=message_key,
           # Не указываем correlation_id и content-type за ненадобностью.
       )
   else:
       error = DataQualityTransferToOpenMetadataError(
           analysis_result=result,
           time=time,
           details=details,
       )
       await dlq_publisher.publish(
           message=generate_message_header(config.ERROR_SCHEMA_ID) + error.SerializeToString(),
           key=message_key,
       )
       logger.warning('Message with the analysis result was sent to DLQ.')


       status = DataQualityTransferToOpenMetadataStatus(
           dq_id=result.dq_id,
           time=time,
           success=False,
           details=details,
       )
       await status_publisher.publish(
           message=generate_message_header(config.STATUS_SCHEMA_ID) + status.SerializeToString(),
           key=message_key,
       )
   logger.info('Operation status was published.')


   await message.ack()

В декораторе помимо топика и консьюмер-группы задаются дополнительные настройки:

  • auto_offset_reset='earliest' позволяет начать обработку с самого раннего сообщения (например, в случае временного простоя консьюмеров);

  • ack_policy=AckPolicy.MANUAL устанавливает ручное управление коммитом, что для нас очень важно;

  • max_poll_interval_ms=10 60 1000 увеличивает стандартный таймаут, так как обработка одного сообщения может затянуться.

Обратите внимание, что фреймворк не поддерживает wire format от Confluent, поэтому приходится работать с заголовком сообщения самостоятельно.

Парсинг результатов

Итак, посмотрим на полученное сообщение DataQualityAnalysisResult.

Сообщение содержит поле table — название проверяемой таблицы. Не зная названия таблицы, мы банально не сможем понять, к какой проиндексированной в OpenMetadata сущности нужно привязать результат анализа. Если же проверка относится к Delta Lake, то вместо table будет указан путь к данным в поле s3_path, из которого можно легко извлечь название сущности. Однако Deequ проверяет любой Spark DataFrame, в том числе полученный из результата произвольного SQL-запроса. В этом случае поля table и s3_path пустые, а в сообщении приходит query с текстом запроса. Чтобы установить FQN сущности в OMD, нам всё равно нужна целевая таблица — значит, ее надо извлечь из запроса. Парсить SQL самостоятельно неприятно, поэтому используем sqlglot. Находим FROM-узел и рекурсивно спускаемся до имени таблицы, попутно проверяя запрос на корректность.

def _extract_table_name(node: sqlglot.expressions.Expression) -> str | None:
    if isinstance(node, sqlglot.expressions.Table):
        schema = node.args.get('db')
        schema_name = schema.args.get('this') if schema else None
        table_name = node.args.get('this').args.get('this')
        return f'{schema_name}.{table_name}' if schema_name else table_name

    if isinstance(node, sqlglot.expressions.Subquery):
        select_node = node.this
        if isinstance(select_node, sqlglot.expressions.Select):
            from_node = select_node.args.get('from_')
            if from_node:
                return _extract_table_name(from_node.this)

    return None

Внутри DataQualityAnalysisResult также находится список проверок в поле constraints, и каждая из них представляет собой результат проверки ограничения Deequ (например, CompletenessConstraint(Completeness(client_id,None))). Именно по этой строке мы должны определить, к какому типу относится ограничение и какой столбец проверяется, если он определен. MetaBridge разбирает строку с помощью простого регулярного выражения.

Типы проверок в Deequ делятся на две группы:

  • табличные (Size, Compliance, Distinctness, Uniqueness) относятся ко всей таблице целиком;

  • столбцовые (Maximum, Minimum, Completeness, MaxLength, MinLength) первым аргументом содержат имя столбца.

При этом название столбца мы вытаскиваем отдельно, потому что в этом случае результат нужно привязать именно к столбцу.

Помните, мы говорили о типах тест-кейсов в OpenMetadata? Так вот, предварительно нам нужно добавить собственные, соответствующие ограничениям Deequ, типы в OMD с помощью запроса POST /v1/dataQuality/testDefinitions. Таким образом мы задали для

  • табличных ограничений — tableSizeCustom, tableComplianceCustom, tableDistinctnessCustom и tableUniquenessCustom;

  • cтолбцовых ограничений — columnMaximumCustom, columnMinimumCustom, columnMaxLengthCustom, columnMinLengthCustom и columnCompletenessCustom.

Такое разделение поможет пользователям каталога легко отличать внутренние типы проверок от наших собственных.

Как результат попадает в OpenMetadata

После определения типа проверки, целевой таблицы, а также столбца, если проверяется именно он, можно приступать к отправке данных в OpenMetadata.

На каждую проверку в рамках конкретного анализа надо выполнить по два запроса. Вначале запросом PUT /v1/dataQuality/testCases мы получаем сущность тест-кейса. Выполняем именно PUT-запрос, чтобы покрыть и создание, и изменение сущности. Затем нужно добавить результат проверки запросом POST /v1/dataQuality/testCases/<uuid>/testCaseResult.

OpenMetadata требует уникальных названий TestCase в рамках одной сущности. То есть, если мы выполняем две проверки типа Completeness для одного и того же столбца, но с разными параметрами фильтрации, нам нужно как-то отразить это в названии тест-кейса. Но что если параметры фильтрации будут очень громоздкими? Тогда и название выйдет очень длинным, что сильно растянет таблицы в интерфейсе и будет, скорее, мешать. Было найдено решение: добавить к имени короткий хеш от полного текста проверки, включающий constraint и SQL-запрос, если именно из него был получен DataFrame.

def generate_test_case_name(name: str, constraint: str, query: str | None = None) -> str:
    raw = f'{constraint}'
    if query:
        raw += f';{query}'
    hashed_name = hashlib.sha256(raw.encode())
    return f'{name} ({hashed_name.hexdigest()[:16]})'

Применяем SHA256 и обрезаем хеш до 16 символов. Так вероятность коллизий пренебрежимо мала, а имя в интерфейсе выглядит читаемо, например Completeness (a1b2c3d4e5f60718). Полный constraint и SQL-запрос при этом не теряются — они добавляются в описание TestCase, так что по нему всегда можно восстановить, что именно проверялось. Да, возможно, вы скажете, что решение с хешем не слишком изящное, но другие нам показались менее удобными.

API OpenMetadata слишком неторопливое

К неспешности API OMD стоит вернуться отдельно, потому что она определила некоторые особенности реализации. Один запрос выполняется в среднем около секунды. Получается, что для результата анализа таблицы с несколькими десятками проверок потребуется около одной минуты. И это при условии, что API не будет отвечать ошибками или сильно тормозить из-за, например, возросшей нагрузки от других операций.

Дело в том, что OpenMetadata написана на Java, у нее богатая модель данных, аудит, валидации на каждом этапе. Требования к задержке у такого сервиса другие. Поэтому мы не боремся со скоростью API, а проектируем систему вокруг этой особенности.

Для Kafka мы установили max.poll.interval.ms равным 10 минутам, так как стандартных пяти иногда не хватает на обработку одного сообщения. При остановке сервиса брокеру дается такой же солидный graceful timeout, чтобы корректно завершить начатые операции. Те же 10 минут прописаны и в terminationGracePeriodSeconds Helm-чарта, чтобы Kubernetes не убил инстанс раньше времени.

Retry и обработка ошибок

Мало того, что OpenMetadata бывает довольно медлительной, так еще и под нагрузкой периодически отвечает ошибками. Этот факт и желание сделать по-настоящему надежный сервис заставили нас сразу задуматься о реализации ключевых механизмов обработок ошибок и нештатных ситуаций.

В используемом нами HTTPX возможности повторных попыток ограничены, поэтому решили воспользоваться tenacity. Получилось достаточно удобно.

@retry(
    retry=retry_if_exception_type(httpx.RequestError),
    wait=wait_random_exponential(multiplier=1, min=1, max=10),
    stop=stop_after_attempt(5),
    reraise=True,
)
async def get_or_create_test_case(client, ...):
    response = await client.put(url, json=test_case_data)
    if response.status_code not in (httpx.codes.OK, httpx.codes.CREATED):
        raise WrongStatusCodeError(...)
    return response

httpx.RequestError — это ошибки уровня соединения, и такие запросы есть смысл повторять, проблема обычно временная. А вот некорректный HTTP-статус превращается в WrongStatusCodeError, который мы прокидываем наверх: если API осмысленно ответил отказом, то повторять тот же запрос, скорее всего, бесполезно.

Отдельно стоит упомянуть статус 409 Conflict при добавлении результата. Это признак того, что такой TestCaseResult уже существует. Поскольку доставка у нас at-least-once, одно и то же сообщение может, теоретически, прийти повторно. И в данном случае ошибка работает как идемпотентность на уровне OMD. В процессе обработки сервис считает дубликаты, а при публикации статуса в соответствующий топик дополнительно указывает, сколько constraint из общего числа уже было загружено ранее.

В планах добавить ещё и circuit breaker. При превышении порога ошибок (скажем, больше половины неудачных запросов за окно) вызывающая сторона сразу получала бы OpenMetadataUnavailable, не делая реального запроса. Это защитило бы и нас от каскадных сбоев, и сам OMD от добивания запросами при критических нагрузках.

Наблюдаемость

Сервис, который разбирает поток сообщений и общается с медлительным внешним API, обязан быть прозрачным. Так что о наблюдаемости мы тоже решили позаботиться.

Логи пишутся в формате logfmt, который стал для наших проектов стандартом. Для просмотра логов используем Grafana.

Любому сервису нужен health check. Очень удобно, что в FastStream есть встроенная поддержка ASGI. Хоть эта поддержка и сильно ограничена, но ее вполне достаточно для реализации обработчиков для запросов GET /internal/alive и GET /internal/ready. Проверка готовности включает тестовый пинг брокера и логирование ошибок:

def readiness(broker: KafkaBroker, broker_logger: logging.Logger) -> ASGIApp:
    healthy_response = AsgiResponse(b'', 204)
    unhealthy_response = AsgiResponse(b'', 500)

    @get async def test(_: Scope) -> AsgiResponse:
        try:
            await broker.ping(timeout=5.0)
        except Exception:
            broker_logger.exception('Kafka not ready.')
            return unhealthy_response
        return healthy_response

    return test

Еще нам бы хотелось собирать метрики в Prometheus. И для этого тоже есть решение в FastStream: с помощью KafkaPrometheusMiddleware мы получаем готовый обработчик GET /metrics с основными показателями.

Результатом довольны

В итоге мы отобразили все результаты проверок качества данных в OpenMetadata. Теперь коллеги могут удобно их просматривать.

Проверки качества данных конкретной таблицы
Проверки качества данных конкретной таблицы
Результат конкретной проверки
Результат конкретной проверки

В описании тест-кейса, как и обещано, указаны ограничение Deequ и SQL-запрос. А вот с графиком вышло не очень красиво. Так как у нас нет числовой метрики для результата проверки, было решено устанавливать значение

  • 1, если проверка завершилась успехом;

  • 0, если проверка была прервана в процессе;

  • -1, если проверки провалилась.

Еще нас попросили сделать небольшой дашборд со статистикой. Он про таблицы только нашего отдела.

Дашборд в Grafana
Дашборд в Grafana

Заключение

Так из примитивного скрипта, читавшего результаты прямо из PostgreSQL, вырос полноценный сервис. Kafka вместо базы как источник данных, Protobuf со схемами в Schema Registry, разбор проверок Deequ и устойчивая к ошибкам отправка в неторопливый API OpenMetadata. Сервис отлично справляется с поставленной задачей.

Сегодня MetaBridge обрабатывает более 20 000 проверок в сутки и покрывает свыше 650 таблиц, а результаты качества данных лежат ровно там, где их ищут пользователи, — рядом с остальными метаданными. В будущем мы ждем кратного увеличения количества проверок, и наш сервис к этому готов.

Останавливаться на качестве данных не хочется. Сейчас мы работаем над тем, чтобы превратить MetaBridge в полноценную систему интеграции инфраструктуры с OpenMetadata.

Главный вывод, который мы для себя сделали: интеграция open source-решения почти никогда не бывает возможна «из коробки». Зато открытый REST API и достаточно гибкая модель данных позволяют встроиться в чужой продукт, не дожидаясь, пока появятся нужные фичи. Надеюсь, наш опыт окажется полезен тем, кто решает похожую задачу.

Комментарии (0)