Назад к блогу

ETLT++: как контракты и мониторинг меняют интеграцию данных

ETLT++: как контракты и мониторинг меняют интеграцию данных

ETL-пайплайны часто дают сбой там, где никто не проверяет корректность данных: формально задача выполнена, а результат уже испорчен. Статья разбирает ETLT++, расширение классической схемы ETL, которое добавляет формальные контракты на данные и наблюдаемость каждого шага — от жизненного цикла запусков в OpenLineage до проверок качества на уровне колонок. Особенно полезно тем, кто строит интеграцию данных и хочет ловить проблемы до того, как они дойдут до потребителей.

Классический ETL описывает поток данных как преобразование «вход → выход», но не фиксирует, что именно произошло внутри и на каких условиях результат считается корректным. ETLT++ добавляет к этому две вещи: формальные контракты на данные и наблюдаемость каждого шага пайплайна. Ниже — механика того, как это устроено в спецификациях и реализациях: жизненный цикл запуска, структура контрактов, привязка проверок качества к датасетам и обработка пограничных случаев.

Жизненный цикл запуска в OpenLineage

Запуск (Run) описывается через события, фиксирующие переход задания в новое состояние. Обязательный минимум — событие START и завершающее событие COMPLETE, FAIL или ABORT; остальные события опциональны. START эмитится в момент старта запуска. Завершающее событие эмитится, когда запуск заканчивается, и его тип отражает исход: COMPLETE, FAIL, ABORT. При старте собираются Run ID, Job ID, тип события START, время события, расположение и версия исходного кода, а также — если уже известны — входы и выходы задания. При завершении собираются Run ID, Job ID, тип COMPLETE, время события и схема выходных наборов данных.

Ключевое свойство — аддитивность метаданных. Если во время выполнения обнаруживаются новые входы или выходы, для них отправляются дополнительные события, но ранее наблюдавшиеся входы и выходы повторно не эмитятся. Это позволяет описывать пайплайн, структура которого раскрывается по мере работы, не переписывая уже зафиксированную историю.

Фасеты: что относится к заданию, а что к датасету

Фасет (Facet) — единица расширения в OpenLineage. Стандартные фасеты разделены на четыре группы.

К заданию относятся Job Facets: место и версия исходного кода, сам исходный код с указанием языка, SQL-запрос (если задание — SQL) и владельцы задания. К датасету — Dataset Facets: схема, экземпляр базы данных, содержащий датасет, состояния жизненного цикла датасета (alter, create, drop, overwrite, rename, truncate), версия при версионировании базой данных, родословная на уровне колонок и владельцы датасета. Отдельно выделены Input Dataset Facets — сюда входит метрика качества данных на уровне датасета и колонок при сканировании: число строк, размер в байтах, число NULL, число уникальных, среднее, минимум, максимум, квантили. Фасеты запуска (номинальное время, родительский запуск, сообщение об ошибке) относятся к запуску, а не к датасету или заданию.

Контракты на входные данные

Media type контракта

В Open Data Contract Standard официальный media type (ранее mime type) задаётся строкой:

application/odcs+yaml;version=3.2.0

Она означает, что контракт представлен в YAML и соответствует версии стандарта v3.2.0 — текущей версии ODCS.

Обязательные поля и валидаторы

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

Проверки полей устроены асимметрично. Непустой runId обязан быть корректным UUID. Если задано имя, но не задано пространство имён, возникает ошибка с сообщением вида «namespace required when name is set». Обе проверки определены в классах, описывающих задания — целевом и исходном, — а не в классах датасетов.

Почему пара namespace+name обязательна целиком

Поля namespace и name по отдельности необязательны, но при установке namespace требуется name, а при установке name требуется namespace. Смысл в том, что пара namespace+name задаёт конкретный внешний источник задания, а их совместное отсутствие означает ссылку на собственное задание события. Поэтому нельзя указать только namespace или только name — это была бы неполная идентификация источника. Аналогичные проверки есть и для целевого задания.

Контракты на выход и промежуточные шаги

Трансформация (LineageTransformation) задаётся типом и необязательными полями: подтип, строковое представление преобразования и признак того, маскирует ли оно исходные данные.

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

Фасет задания содержит список записей, описывающих целевые сущности и питающие их источники. Запись целевого задания имеет тип JOB, необязательные namespace и name (при совместном отсутствии цель — собственное задание события), runId (когда связь привязана к конкретному исполнению) и список источников. К runId применяется та же проверка на UUID, к name — та же проверка парности с namespace. Те же валидаторы определены и для исходного задания.

Проверки качества: уровни и ассерты

Ассерты датасета

Ассерт (Assertion) описывается фасетом входного датасета, который содержит список выполненных тестов и их результаты. У каждого ассерта обязательны классификация теста и признак успеха. Классификация теста указывает, что именно проверяется: not_null — отсутствие пустых значений, unique — уникальность значений, row_count — число строк, freshness — свежесть данных, custom_sql — произвольный SQL-запрос. Поле успеха принимает значения true (проблем не найдено) и false (проблемы найдены) независимо от уровня серьёзности: тест может провалиться, не блокируя пайплайн, если серьёзность — warn. Поле column указывает колонку, к которой относится тест, и при пустом значении тест относится ко всему датасету. Поле severity задаёт уровень серьёзности: обычно error (провал блокирует пайплайн) или warn (провал даёт только предупреждение). Необязательные поля включают имя теста, человекочитаемое описание, ожидаемое и фактическое значения (сериализованные строками), тело проверки, формат этого тела (например, sql, json, expression) и произвольные параметры. Привязка к датасету обеспечивается тем, что фасет является фасетом входного датасета, а его схема возвращается как URL спецификации DataQualityAssertionsDatasetFacet.

Сравнение счётчиков двух таблиц

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

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

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

Сравнение с ожидаемым числом строк

Проверка равенства числа строк заданному значению берёт ожидаемое из параметров успеха по ключу value, наблюдаемое — из метрики числа строк, а успех определяет их равенством. Параметры метода — метрики, конфигурация времени выполнения и движок исполнения, но используется только словарь метрик. Ожидаемое значение задаётся полем value («ожидаемое число строк»), метрика-зависимость — число строк.

При рендеринге описания создаётся конфигурация рендерера из конфигурации, результата и конфигурации времени выполнения; через подстановку недостающих значений в шаблон «Must have exactly $value rows.» подставляется значение, а стилизация берётся из конфигурации времени выполнения. Возвращается список отрендеренного строкового шаблона с типом блока string_template, содержащим шаблон, параметры и стилизацию.

Мониторинг: метрики и счётчики

Парсер проверки числа строк

Парсер проверки для типа row_count объявляет это имя типа и создаёт реализацию проверки из контракта, колонки и YAML-описания. В конструкторе порог создаётся со значением по умолчанию: тип SINGLE_COMPARATOR с требованием «больше нуля». Имя метрики вычисляется из имени типа проверки и колонки, а сводка строится из порога, если он есть, иначе — как строка с пометкой о недействительном пороге.

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

if self.check_filter:
    return SUM(CASE_WHEN(SqlExpressionStr(self.check_filter), LITERAL(1)))
else:
    return COUNT(STAR())

Преобразование значения из базы возвращает целое, если значение не None, иначе 0 — потому что выражение вида SUM(CASE WHEN "id" IS NULL THEN 1 ELSE 0 END) даёт NULL при отсутствии строк. Настройка метрик создаёт метрику подсчёта строк, а оценка получает её значение и передаёт в оценку порога.

Почему SUM(CASE WHEN ...) может вернуть NULL

При наличии фильтра SQL-выражение возвращает сумму единиц по условию. Как отмечено в комментарии, такое выражение даёт NULL, если строк нет. Преобразование значения из базы обрабатывает это: возвращает целое при непустом значении, иначе 0. Настройка метрик создаёт метрику подсчёта строк, а оценка получает значение метрики и передаёт его в оценку порога.

Мониторинг: свежесть и прогнозы

Проверка соответствия числа строк в батче прогнозу Prophet-модели на заданную дату — это BatchExpectation. Её метод оценки принимает метрики, конфигурацию времени выполнения и движок исполнения. Из метрик берётся число строк батча, из конфигурации — сериализованная модель и дата. Модель восстанавливается десериализатором, после чего прогноз строится вызовом predict на DataFrame с колонкой ds, содержащей дату. Из прогноза берутся точечное значение и нижняя и верхняя границы.

Аномалией считается ситуация, когда число строк не попадает строго внутрь интервала:

in_bounds = (forecast_lower_bound < batch_row_count) & (
    batch_row_count < forecast_upper_bound
)

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

Реакция на сбои и владельцы

Связь результата проверки с конкретным датасетом обеспечивается тем, что фасет ассертов является фасетом входного датасета и содержит список проверок с их результатами. Каждая проверка указывает колонку (или её отсутствие означает весь датасет) и признак наличия проблем.

Привязка результата проверки к владельцу источника в спецификации не описана. Владение упоминается только как фасет задания («владельцы задания») и фасет датасета («владельцы датасета»), без связи с результатами проверок. Фасеты выходного датасета в документе включают только статистику вывода — размер записанных в датасет данных (число строк и размер в байтах).

Пограничные случаи и компромиссы

Для полей namespace и name действуют две симметричные проверки. Если задано name, но отсутствует namespace, возбуждается ValueError с сообщением «namespace required when name is set». Если задано namespace, но отсутствует name, возбуждается ValueError с сообщением «name required when namespace is set». Для runId определён валидатор, который при непустом значении пытается преобразовать его в UUID. Поля namespace и name необязательны, и при совместном отсутствии означают ссылку на собственное задание события. Что именно видит потребитель метаданных при частичной эмиссии событий, в спецификации не описано.

Что из этого следует на практике

  • Обязательный минимум событий — START и одно из COMPLETE/FAIL/ABORT. Всё остальное опционально, а метаданные аддитивны: новые входы и выходы описываются дополнительными событиями без повторной эмиссии уже наблюдённых.
  • Контракт на датасет неполон без пары namespace+name: указать только одно из полей нельзя — валидатор немедленно возбудит ошибку. Совместное отсутствие пары — это осознанная ссылка на собственное задание события, а не пропуск данных.
  • Непустой runId обязан быть корректным UUID; иначе валидатор упадёт ещё до отправки события.
  • Привязка проверок качества к датасету идёт через фасет входного датасета, а не выходного. Фасеты выходного датасета ограничены статистикой вывода, поэтому связать результат проверки с владельцем источника через них нельзя.
  • Уровень серьёзности (error или warn) не влияет на поле успеха: тест может провалиться, не блокируя пайплайн. Блокировка определяется именно серьёзностью, а не фактом провала.
  • Сравнение числа строк двух таблиц — единственная проверка, вычисляющая одну метрику по двум доменам, и реализована она ручным копированием конфигурации метрики с переименованием обеих. Это накладывает порядок: обе метрики вычисляются до сравнения.
  • При подсчёте строк с фильтром SQL-выражение может вернуть NULL при отсутствии строк; преобразование значения из базы обязательно приводит NULL к нулю, иначе пороговая проверка получила бы некорректное значение.
  • Проверка по Prophet-модели считает аномалией строгое непопадание внутрь интервала: значения ровно на границе уже аномалия.

Где смотреть в коде

Источники

Похожее