АрхитектураАрхитектура данныхИнженер по платформе данных

Как предотвратить попадание в аналитическое хранилище событий, несовместимых с согласованной схемой?

Как предотвратить попадание в аналитическое хранилище событий, несовместимых с согласованной схемой?

Проходите собеседования с ИИ помощником Hintsage

Краткий ответ

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

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

Исторический контекст

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

Реестр схем и формальные правила совместимости появились как способ сделать контракт машинно проверяемым. Вместо договорённости «не ломать старых потребителей» система может автоматически проверять изменения схемы и отклонять несовместимые версии до их публикации.

Постановка проблемы

Событие может нарушать контракт несколькими способами: отсутствует обязательное поле, тип значения изменён, поле переименовано без сохранения обратной совместимости или структура вложенного объекта стала другой. Если такое сообщение сразу попадает в аналитическое хранилище, ошибка проявляется поздно — во время загрузки, построения витрины или выполнения отчёта.

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

Подробное решение

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

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

Важно различать несколько уровней контроля:

  • структурная валидация проверяет форму сообщения и типы;
  • совместимость схем определяет, смогут ли старые потребители читать новые сообщения;
  • семантическая валидация проверяет бизнес-смысл, например положительность суммы или существование идентификатора;
  • контроль качества выявляет выбросы, пропуски и неожиданные распределения уже на уровне данных.

Для эволюции схемы обычно безопаснее добавлять необязательное поле, чем переименовывать или удалять существующее. Но допустимый вариант зависит от выбранного режима совместимости: обратная совместимость ориентирована на чтение новых данных старым потребителем, прямая — на чтение старых данных новым потребителем, а строгий режим требует выполнения обоих направлений.

Проверка только у производителя недостаточна: ошибочный или устаревший клиент может обойти это правило. Проверка только в аналитическом хранилище также запаздывает и может позволить некорректным данным пройти через промежуточные системы, поэтому защитные проверки разумно размещать как минимум у производителя и на входе платформы данных.

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

Ситуация из практики

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

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

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

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

Что кандидаты часто упускают

1. Достаточно ли зарегистрировать схему, чтобы гарантировать качество данных?

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

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

2. Где лучше отклонять сообщение — у производителя или у потребителя?

Проверка у производителя дешевле с точки зрения предотвращения распространения ошибки и быстрее возвращает обратную связь владельцу источника. Но полностью полагаться на неё нельзя: существуют старые версии сервисов, внешние источники и ошибки маршрутизации.

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

3. Почему нельзя разрешить любое изменение схемы, если потребители используют только нужные им поля?

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

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