Что такое Enrichment-архитектура (архитектура обогащения данных)
Подборки сервисов

Что такое Enrichment-архитектура (архитектура обогащения данных)

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

Елена Кравцова
Елена Кравцова
Редактор и автор статей23 мин

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

Термин data enrichment, «обогащение данных», устоялся в data engineering. Сочетание Enrichment architecture употребляют как название паттерна или контура решения, но единого стандарта с обязательным набором технологий у него нет. Поэтому полезнее определять такую архитектуру через свойства: исходные данные сохраняются, добавленные факты имеют источник и версию, сбои не теряют записи, а результат можно проверить, переобработать и откатить.

Что именно считается обогащением данных

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

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

Практический тест прост: после операции потребитель знает о событии или сущности что-то полезное, чего не было в исходной записи. Замена RU на RUS меняет формат. Добавление к country_code=RU налогового региона, часового пояса и разрешённых способов доставки создаёт новый контекст.

Границы с ETL, нормализацией, RAG и feature engineering

Соседние термины пересекаются, но отвечают на разные архитектурные вопросы.

Понятие Главный вопрос Как связано с обогащением
Ingestion Как доставить данные из источника? Доставка без новых сведений остаётся ingestion. Она может запускать конвейер обогащения
Transformation Как изменить или вычислить данные? Более широкая категория. Обогащение часто реализовано как transform, который добавляет контекст
Нормализация Как привести значения к единому формату? Обычно предшествует сопоставлению: телефон приводят к E.164, валюту — к ISO-коду, адрес разбирают на поля
Data augmentation Как расширить набор данных? В машинном обучении часто означает создание синтетических обучающих примеров; enrichment обычно дополняет существующую запись
ETL/ELT (extract, transform, load / extract, load, transform) В каком порядке извлечь, преобразовать и загрузить? Обогащение может быть стадией ETL до загрузки или ELT после неё
CDC (Change Data Capture, захват изменений данных) Как передавать изменения из рабочей базы? CDC поддерживает свежую копию справочника или запускает перерасчёт, но не добавляет смысл сам по себе
Reverse ETL Как вернуть подготовленные данные из хранилища в рабочие системы? Часто доставляет уже обогащённые сегменты и показатели в CRM, рекламу или приложение
Entity resolution (сопоставление сущностей) и MDM (Master Data Management, управление мастер-данными) Какие записи описывают одну сущность и какой вариант эталонный? Дают устойчивый entity_id и golden record (эталонную запись), к которым затем присоединяют атрибуты
Feature engineering (проектирование признаков) Какие признаки нужны модели? Специализированное обогащение для машинного обучения (ML). Для него важна одинаковая логика обучения и online-инференса
RAG (Retrieval-Augmented Generation, генерация с дополненным поиском) Какой контекст дать языковой модели для текущего ответа? RAG обычно добавляет контекст к запросу, не изменяя исходную запись. Подготовка документов, тегов и embeddings для индекса может быть enrichment

Эта граница влияет на ответственность. ETL-оркестратор отвечает за запуск задач, MDM — за эталонную сущность, feature store — за согласованную выдачу признаков, reverse ETL — за запись в целевой сервис. Enrichment-контур связывает эти части только там, где реально добавляется и проверяется новый атрибут.

Как выглядит производственный контур

Базовая схема отделяет неизменённый вход, справочный контекст и версионированный результат.

                                  +--------------------+
Системы-источники -> raw storage -> контракт и очередь | -> оркестратор
                                  +--------------------+        |
                                                               +-- детерминированные правила
Reference DB -> CDC -> reference topic/state/cache ------------+-- сторонний API
                                                               +-- языковая / ML-модель
                                                               `-- ручная проверка
                                                                        |
                         +----------------------------------------------+
                         v
              quality gate -> enriched topic / side table -> serving
                    |                    |                    |
                    |                    |                    `-> API, BI, ML,
                    |                    |                        reverse ETL
                    |                    `-> lineage, версии, аудит
                    `-> quarantine / retry / DLQ -> replay и backfill

В этой схеме нет обязательного Kafka, Redis или конкретного облака. Есть обязанности, которые кто-то должен выполнить.

Источник и неизменяемый raw-слой

Raw-слой хранит событие в том виде, в котором система его приняла, вместе с идентификатором, временем события и версией схемы. Он нужен для расследований и replay. Если новый обогатитель оказался ошибочным, команда повторно запускает данные из raw, а не пытается восстановить прежние значения из перезаписанной таблицы.

Raw не обязательно означает бесконечное хранение. Срок определяют бизнес-цель, стоимость и требования к персональным данным. Важно, чтобы удаление по запросу субъекта данных распространялось на raw, кэш, результаты, DLQ и резервные копии в соответствии с применимыми правилами.

Контракт и схема

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

Контракт также отвечает на вопросы:

  • обязательный ли атрибут и можно ли выдать частичный результат;
  • чем unknown отличается от not_applicable и ошибки;
  • какая команда владеет полем;
  • разрешено ли передавать поле внешнему поставщику;
  • какие изменения схемы обратимо совместимы;
  • когда нарушение нужно отклонить, поместить в quarantine или отправить в dead letter queue (DLQ).

Реестр схем помогает технически проверять совместимость. Полноценные data contracts добавляют ограничения значений, метаданные, теги персональных данных, правила миграции и политики обработки.

Журнал или очередь

Журнал событий сохраняет порядок внутри partition и позволяет нескольким потребителям независимо читать историю. Очередь распределяет задания между workers. Kafka уместен, если нужны replay, несколько потоковых потребителей, большой throughput и длительно живущий log. SQS, RabbitMQ, NATS, облачный Pub/Sub или таблица заданий в PostgreSQL могут быть проще, когда поток меньше и потребитель один.

Выбор транспорта не отменяет дубликаты. Producer может не получить подтверждение и отправить событие снова, consumer — упасть после внешнего вызова, но до фиксации offset. Поэтому практическая гарантия строится вокруг идемпотентного эффекта, а не обещания «exactly once» на слайде.

Оркестратор и независимые обогатители

Оркестратор знает граф зависимостей и политику завершения. География и категория товара могут вычисляться параллельно, а риск — только после entity resolution и географии. Каждый обогатитель (enricher) получает версионированный вход и возвращает структурированный результат со статусом.

Разделение enrichers уменьшает радиус сбоя. Изменение классификатора обращений не должно заставлять повторно покупать геоданные. Для каждой стадии нужны собственные timeout, quota, retry policy, версия и целевой уровень сервиса (Service Level Objective, SLO).

Кэш, state store и справочные данные

Одинаковые ключи не стоит отправлять во внешний API миллионы раз. Точный кэш по ключу сокращает задержку и стоимость. Time to live (TTL) ограничивает максимальное окно устаревания, но не доказывает свежесть: значение могло измениться через секунду после заполнения кэша.

Stateful stream processor держит локальное состояние рядом с вычислением. Справочные изменения поступают через CDC, а событие присоединяется к текущей или исторической версии. Это быстрее удалённого запроса на каждую запись, но добавляет задачи начального snapshot, восстановления state, rekey, контроля размера и временной корректности.

Идемпотентность, дедупликация и фиксация

Один логический вход должен иметь устойчивый event_id. Результат удобно идентифицировать составным ключом:

(event_id, enricher_name, enricher_version, input_hash)

Повтор с тем же ключом возвращает или обновляет тот же логический результат. Новая версия алгоритма создаёт новую запись. В PostgreSQL это можно закрепить уникальным индексом и INSERT ... ON CONFLICT. Внешний API требует собственного idempotency key, если поставщик его поддерживает. Если не поддерживает, перед вызовом и после него нужен журнал попыток, а для необратимых операций — особая осторожность.

Retry, circuit breaker и DLQ

Повторять стоит только временные ошибки: timeout, разрыв соединения, 429 и часть 5xx. Задержка растёт экспоненциально и получает случайный jitter, чтобы тысячи workers не атаковали восстановившийся сервис одновременно. Число попыток и общий retry budget ограничивают.

Ошибку контракта, запрещённый доступ или неверный API-ключ повтор обычно не исправит. Такие записи получают reason code и уходят в quarantine либо DLQ. У DLQ должны быть владелец, срок хранения, оповещение, инструмент безопасного replay и метрика возраста самой старой записи. Иначе она превращается в скрытое хранилище потерянных данных.

Circuit breaker временно прекращает запросы к устойчиво неработающему поставщику. Дальше действует заранее выбранная деградация: вернуть pending, выдать допустимое stale-значение, пропустить необязательное поле или остановить публикацию обязательного результата. Подставлять правдоподобное значение нельзя.

Provenance, confidence и версии

Потребителю недостаточно увидеть country=RU. Нужны сведения, откуда оно появилось и к какому времени относится:

  • имя и версия enricher;
  • версия набора данных или API;
  • source_as_of, event_time и processed_at;
  • хеш входа;
  • evidence или matched key;
  • статус качества и проверки;
  • измеренная задержка и стоимость;
  • модель и версия prompt для большой языковой модели (Large Language Model, LLM);
  • автор и причина ручного решения.

Confidence полезен только при понятной калибровке. Число 0.97, которое языковая модель написала о собственном ответе, не становится вероятностью точности. Для правил можно хранить однозначный match; для ML — откалиброванный score и версию модели; для человека — статус review и reason code.

Версия результата не должна уничтожать предыдущую. Удобна модель append-only результатов плюс указатель current. Тогда rollback меняет указатель или serving policy, а история остаётся для расследования.

Наблюдаемость, quality gate и serving

Quality gate решает, можно ли публиковать запись. Условия могут включать обязательные поля, диапазоны, свежесть источника, отсутствие конфликтов и минимальный проверенный score. Статусы pending, partial, ready, rejected и needs_review делают eventual consistency видимой для потребителя.

Serving-слой выдаёт результат через API, таблицу, поисковый индекс, feature store или downstream topic. Reverse ETL может перенести сегмент или показатель в CRM. Запись во внешний SaaS требует особенно строгого контроля: целевая схема принадлежит поставщику, rate limits отличаются, а массовый rollback иногда невозможен.

Где выполнять enrichment

Место обогащения выбирают по требуемой задержке, свежести, цене и допустимой связанности с поставщиком.

Режим Когда подходит Сильная сторона Основной риск
Request-time lookup Низкое число запросов в секунду (RPS), нужен самый свежий ответ, допустима зависимость от API Нет предварительной материализации Внешняя задержка и отказ попадают в пользовательский запрос
Write-time synchronous Без обязательных полей запись нельзя публиковать Downstream получает полный контракт Замедляет ingestion, расширяет радиус сбоя
Write-time asynchronous Высокий поток, enrichers независимы, допустима eventual consistency Буферизация, масштабирование и изоляция отказов Потребитель должен понимать pending/partial
Read-time join Логика часто меняется, запросов мало, важна свежесть справочника Не создаёт множество материализованных копий Повторные вычисления, сложная историческая воспроизводимость
Batch materialization Аналитика, ночной пересчёт, большая история, свежесть в часах/сутках Дешёвые bulk-операции и простой backfill Результат устаревает между запусками
Streaming stateful join Сотни и тысячи событий в секунду, reference data доступна через CDC Низкая задержка без сетевого lookup на каждую запись Управление state, ключами, checkpoint и временем
Human-in-the-loop Ошибка дорога, совпадение неоднозначно, автоматике нельзя доверить решение Контролируемое качество и объяснимость Очередь, стоимость, SLA и согласованность экспертов

Batch или streaming

Batch работает с ограниченным набором: дневной partition, выгрузка клиентов, весь каталог. Он экономичен при bulk API и удобен для тяжёлых join. Streaming обрабатывает события по мере поступления и нужен для antifraud, персонализации или оперативного маршрута.

Гибрид встречается чаще чистых крайностей. Live-поток использует быстрый кэш и ограниченный набор полей, а ночной batch сверяет полноту и исправляет накопленные расхождения. Backfill работает отдельной очередью с меньшим приоритетом, чтобы история не вытеснила живые события.

Синхронно или асинхронно

Синхронный вызов удерживает worker до ответа. Если API отвечает 250 мс, один worker теоретически делает около четырёх последовательных вызовов в секунду. Асинхронный клиент держит несколько запросов in flight и использует время ожидания для другой работы. Apache Flink Async I/O также учитывает порядок, event time, checkpoint и повтор in-flight запросов после восстановления.

Асинхронность не означает бесконечную concurrency. Её ограничивают rate limit поставщика, пул соединений, память, максимальная очередь и backpressure. Ordered-выдача сохраняет порядок, но медленный ответ задерживает последующие. Unordered-выдача повышает throughput, если downstream не зависит от порядка.

Write-time или read-time

Write-time материализация ускоряет многократное чтение и фиксирует контекст, известный в момент события. Цена — хранение копий и необходимость переобогащения при исправлении справочника.

Read-time join применяет актуальный контекст при каждом запросе. Для отчёта «как мы знали тогда» этого мало: понадобится историческая версия справочника и temporal join. Полезно хранить как минимум три времени:

  • event_time — когда произошло событие;
  • source_as_of — на какой момент действовал присоединённый факт;
  • processed_at — когда система вычислила результат.

Для ML добавляется created_at: поздняя корректировка с прошлым event_time могла ещё не существовать в момент прогноза. Point-in-time join должен исключать утечку будущего знания в обучающую выборку.

Откуда брать дополнительные данные

Reference data и локальный join

Страны, валюты, категории, тарифы и характеристики товаров удобно хранить как версионированные справочники. Маленький и редко меняющийся набор можно загрузить в память каждого worker. Большой или часто меняющийся — передавать через CDC в state store либо запрашивать temporal lookup.

При выборе проверяют:

  • размер полного набора и рабочей части;
  • частоту изменений;
  • ключ и кардинальность;
  • допустимое окно устаревания;
  • нужна ли историческая версия;
  • поведение при отсутствующем и неоднозначном ключе.

Сторонние API

API подходит для адресов, компаний, риска, геоданных и других лицензируемых наборов. Вместе с ответом система получает эксплуатационные ограничения: цену вызова, quota, rate limit, территорию обработки, срок хранения и условия производного использования.

Перед интеграцией нужен контракт деградации. Если геокодер недоступен, можно ли принять заказ без координат? Если risk-provider не ответил, следует ли остановить операцию? Ответ зависит от бизнес-риска, а не от удобства разработчика.

Секреты хранят в secret manager и не добавляют в event payload, логи или DLQ. Доступ дают минимальной service identity. Ответ provider проходит проверку схемы так же строго, как внутреннее событие.

Entity resolution и MDM

До присоединения профиля нужно понять, к какой сущности относится запись. Email, телефон и имя могут конфликтовать; один человек может иметь несколько аккаунтов, а общий телефон — принадлежать семье. Entity resolution использует правила, вероятностное сопоставление или ML и выдаёт устойчивый match ID.

Master Data Management (MDM) идёт дальше: стандартизирует, сопоставляет, объединяет и валидирует записи, формируя golden record. Низкая уверенность или нарушение правил отправляются data steward на разбор. Enrichment затем присоединяет атрибуты к устойчивому ID. Иногда результат сопоставления и golden record сами служат обогащёнными данными.

LLM-enrichment

Языковая модель умеет извлекать сущности из письма, классифицировать тему, составлять краткое резюме, предлагать теги и приводить неструктурированный текст к JSON. Это полезно там, где жёсткие правила слишком хрупки.

Надёжный LLM-enricher ограничивает свободу модели:

  1. Вход проходит redaction и проверку разрешённой цели обработки.
  2. Prompt требует структурированный output по версии схемы.
  3. Для извлечённых фактов модель возвращает evidence spans из исходного текста.
  4. Валидатор проверяет enum, диапазоны, обязательные поля и логические связи.
  5. Детерминированные справочники подтверждают то, что можно подтвердить без модели.
  6. Неуверенный или рискованный результат получает needs_review, а не догадку.
  7. Система хранит model, prompt version, input hash, output и решение reviewer.

JSON, прошедший schema validation, может содержать вымышленный факт. Поэтому hallucination — отдельный риск качества. Для финансовых, медицинских, юридических и иных значимых решений нужен более строгий review и контроль применимых правил об автоматизированных решениях.

Human review

Человек нужен для неоднозначных сущностей, новых категорий и дорогих ошибок. Очередь должна показывать исходную запись, версии источников, варианты и evidence. Reviewer выбирает результат и reason code, а система пишет аудит.

Полезно измерять долю ручных исправлений по версии enricher и согласие независимых reviewers. Если override rate растёт, проблема может быть в drift, новом типе входа или плохом справочнике, а не в «невнимательности операторов».

Контракт события и результата

Рассмотрим синтетическое событие заказа. Адрес из документационного диапазона и результат ниже служат только примером структуры, а не реальным GeoIP-ответом. В raw нет региона, сегмента и риска.

{
  "event_id": "01K4GQ8N4Y7AX2Z2S45R6TQJPG",
  "event_type": "order.created",
  "schema_version": 2,
  "event_time": "2026-09-06T10:15:23Z",
  "observed_at": "2026-09-06T10:15:24Z",
  "entity": {
    "customer_id": "c_4821",
    "product_id": "p_917"
  },
  "payload": {
    "amount": 12990,
    "currency": "RUB",
    "ip": "203.0.113.8"
  },
  "policy": {
    "purpose": "fraud_prevention",
    "region": "RU"
  },
  "trace_id": "8f34f42e6e9c4abc"
}

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

{
  "event_id": "01K4GQ8N4Y7AX2Z2S45R6TQJPG",
  "enricher": "geoip",
  "enricher_version": "geoip-v4",
  "input_hash": "sha256:79d3...",
  "source": {
    "dataset": "geo-db",
    "version": "2026-09-05",
    "as_of": "2026-09-05T00:00:00Z"
  },
  "status": "succeeded",
  "value": {
    "country": "RU",
    "city": "Москва"
  },
  "confidence": null,
  "evidence": {
    "matched_prefix": "203.0.113.0/24"
  },
  "processed_at": "2026-09-06T10:15:24.180Z",
  "latency_ms": 8,
  "cost_units": 0
}

confidence здесь null: точный lookup вернул совпадение по версии базы, и придумывать score не требуется. У вероятностного entity matching поле имело бы значение вместе с названием модели и способом калибровки.

Материализованное представление для serving может выглядеть так:

{
  "event_id": "01K4GQ8N4Y7AX2Z2S45R6TQJPG",
  "status": "partial",
  "source_event_version": 2,
  "enrichment_set_version": "order-context-v7",
  "attributes": {
    "geo": {"country": "RU", "city": "Москва"},
    "customer_segment": "returning",
    "risk_band": null
  },
  "pending": ["risk-score-v3"],
  "ready_at": null,
  "updated_at": "2026-09-06T10:15:24.230Z"
}

Потребитель видит, что география и сегмент готовы, а риск ещё считается. Если риск обязателен, quality gate не переведёт запись в ready.

Как считать throughput, кэш и стоимость

Начинать нужно не с числа workers, а с входного потока, пика, задержки зависимостей и внешних лимитов.

Предположим, система принимает 10 миллионов событий в сутки:

средний поток = 10 000 000 / 86 400 = 115,7 события/с
пиковый поток при коэффициенте 5 = 579 событий/с

Если внешний API отвечает за 250 мс на 95-м процентиле задержки (P95), минимальное число одновременных запросов по закону Литтла:

in-flight = 579 * 0,25 = 144,7
с запасом 30%: 144,7 * 1,3 ≈ 188 запросов

Это расчёт ёмкости, а не разрешение открыть 188 соединений. Если provider допускает только 100 запросов в секунду, система должна применить cache, batching, предварительное обогащение справочника, более высокий тариф или отложенную очередь. Простое увеличение concurrency создаст 429 и retry storm.

Теперь оценим кэш. При hit rate 95% наружу уйдёт 5% событий:

вызовы API = 10 000 000 * 0,05 = 500 000 в сутки

При условной цене $0,002 за вызов:

Вариант Вызовы в сутки Условная стоимость
Без кэша 10 000 000 $20 000
Hit rate 95% 500 000 $1 000
Разница 9 500 000 $19 000

Цифры иллюстрируют формулу, а не тариф конкретного сервиса. Hit rate зависит от повторяемости ключей, кардинальности, TTL и распределения нагрузки. Длинный TTL повышает hit rate, но увеличивает вероятность stale-ответа. Инвалидация по CDC может дать и высокую долю попаданий, и приемлемую свежесть, если источник передаёт изменения надёжно.

Основные метрики кэша:

hit_rate = hits / (hits + misses)
origin_qps = misses / second
stale_age = now - source_as_of

Нужно также видеть evictions, hot keys, stampede после истечения TTL и долю negative-cache записей. Кэшировать «не найдено» полезно коротко: иначе новый объект останется невидимым до длинного TTL.

Главные компромиссы

Задержка против полноты

Параллельные enrichers сокращают суммарное время, но обязательная самая медленная стадия определяет latency. Можно выдать partial и досчитать остальное позже либо ждать полный результат. Выбор фиксируют в контракте, иначе разные потребители начнут трактовать null по-разному.

Свежесть против стоимости

Request-time API даёт свежий ответ ценой каждого чтения. Материализация и cache уменьшают расходы, но добавляют staleness. CDC-инвалидация улучшает баланс, однако сама становится критической зависимостью: lag и пропущенные изменения нужно измерять.

Согласованность против доступности

Синхронный gate может отказаться публиковать заказ без risk-score. Это повышает целостность, но отказ risk-provider остановит бизнес-путь. Async-подход сохранит заказ и поставит риск в pending, зато downstream должен уметь ждать и не совершать преждевременное действие.

Скорость против воспроизводимости

Join с «текущим профилем» прост и быстр. Он не отвечает, какой профиль был известен в момент старого события. Для аудита, ML и финансовых расчётов нужны versioned reference data, event time и point-in-time join.

Универсальность против управляемости

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

Защита данных, право и лицензии

Обогащение увеличивает ценность данных и одновременно повышает риск. Набор «email + место + интересы + доход + риск» чувствительнее каждой части по отдельности.

До подключения источника команда должна ответить:

  • для какой цели нужно каждое поле и есть ли правовое основание;
  • можно ли отправлять исходные поля поставщику и в какой юрисдикции он их обработает;
  • как выполнить минимизацию, удаление, исправление и ограничение срока хранения;
  • допускает ли лицензия хранение, производные признаки, передачу дочерним системам и ML-использование;
  • можно ли объяснить и оспорить значимое автоматизированное решение;
  • как удалить данные из raw, cache, state, результатов, логов, DLQ и backfill-копий;
  • какие требования 152-ФЗ, GDPR и отраслевых норм применимы к конкретному проекту.

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

Технические меры включают field-level allowlist, tokenization, шифрование, раздельные ключи, минимальные service roles, redaction логов и сроки хранения по типу данных. Provenance позволяет найти все результаты конкретного поставщика и удалить либо переобогатить их при отзыве лицензии.

Как переживать отказ поставщика и ошибку алгоритма

Для каждого enricher заранее задают таблицу поведения.

Ситуация Retry Действие
Timeout, 408, временный 5xx Ограниченный, backoff + jitter Повторить в пределах budget, затем DLQ или pending
429 После Retry-After, с глобальным rate limiter Снизить concurrency, не создавать retry storm
400: неверная схема Нет Quarantine, alert владельцу контракта
401/403 Нет до исправления доступа Открыть circuit, alert, не писать секрет в лог
Ответ неполный По политике конкретного API Принять partial, запросить альтернативу или review
Cache недоступен Иногда Ограничить fallback на origin, чтобы не обрушить provider
LLM нарушил JSON schema Ограниченный Повтор с тем же versioned prompt или needs_review
LLM дал неподтверждённый факт Автоповтор не помогает Evidence check, deterministic validation, human review
Ошибка новой версии enricher Нет для той же версии Остановить rollout, вернуть current pointer, replay на старой версии

Rollback работает только при сохранённом raw, versioned output и управляемом serving pointer. Если обогащение без истории перезаписало профиль, прежнее значение уже не восстановить. Если reverse ETL разослал ошибку по CRM и рекламным кабинетам, понадобится компенсирующая синхронизация, а некоторые внешние действия окажутся необратимыми. Поэтому новые версии сначала считают в shadow mode.

Варианты реализации

PostgreSQL, Redis и workers

Такой стек подходит для первой версии и умеренной нагрузки.

API/импорт
   |
   v
PostgreSQL: raw_events + enrichment_jobs + enrichment_results
   |                   |
outbox              workers -- Redis cache/rate limiter --> external API
   |                   |
broker/webhook          `-> review_queue / dead_letters
   |
serving table or API

В PostgreSQL можно завести:

CREATE TABLE enrichment_results (
  event_id uuid NOT NULL,
  enricher text NOT NULL,
  enricher_version text NOT NULL,
  input_hash text NOT NULL,
  status text NOT NULL,
  value jsonb,
  source_as_of timestamptz,
  processed_at timestamptz NOT NULL DEFAULT now(),
  PRIMARY KEY (event_id, enricher, enricher_version, input_hash)
);

Workers забирают задания короткими транзакциями через FOR UPDATE SKIP LOCKED, обрабатывают их вне транзакции и фиксируют результат идемпотентным upsert. LISTEN/NOTIFY можно использовать для пробуждения, но durable truth остаётся в таблице. Transactional outbox связывает изменение БД и публикацию события без опасного dual write.

Redis подходит для cache-aside, rate limit и коротких coordination locks. Он не должен становиться единственным хранилищем raw, DLQ или аудита. При отказе Redis нужна защита origin: semaphore, stale-if-error и ограниченный fallback. Читать полный обзор сервиса Redis.

Kafka, Flink и локальное состояние

Этот вариант нужен при большом потоке, нескольких потребителях и temporal joins.

PostgreSQL -- Debezium CDC --> customer-reference topic --+
                                                          |
orders topic --> rekey --> Flink stream-table join <------+--> enriched-orders
                                |
                                +--> async API lookup --> retry/DLQ topics
                                `--> checkpoints/state

Ключ события должен совпадать с ключом справочника. CDC сначала делает согласованный snapshot, затем передаёт committed changes. Flink materializes reference state и присоединяет его локально. Для большого внешнего справочника применяется lookup join или async I/O.

Нужно решить, какое состояние корректно: текущее на processing time, действовавшее на event time или известное системе на тот момент. Последний вариант требует учитывать и время бизнес-факта, и время его появления.

Checkpoint восстанавливает оператор и in-flight записи, но внешний API может увидеть повтор. Поэтому sink и side effects всё равно делают идемпотентными.

Serverless

Serverless-контур можно собрать из managed queue/stream, фильтра, функции, workflow и key-value store. Например:

SQS/Kinesis -> EventBridge Pipes -> Lambda/Step Functions -> DynamoDB
                    |                         |
                  filter                  retry/DLQ
                    `-> enrichment call -> EventBridge target

Аналоги есть в Google Cloud и Azure: Pub/Sub или Event Hubs, Functions/Cloud Run, Workflows/Durable Functions, managed state и dead-letter queue.

Serverless удобен для всплесков, нерегулярной нагрузки и небольшой операционной команды. Проверить нужно максимальную длительность, cold start, лимиты concurrency, стоимость сетевых вызовов, порядок событий, размер batch и удобство массового replay. Сложный граф лучше вынести в workflow engine, а не наращивать цепочку скрытых триггеров.

Тестирование и миграция

Какие тесты нужны

Unit-тест счастливого пути не проверяет архитектуру. Нужны следующие уровни.

  • Контрактные тесты: совместимость схем, обязательные поля, enum, PII tags, старый producer с новым consumer.
  • Golden fixtures: известный вход и ожидаемый результат для каждой версии enricher.
  • Property-based тесты: диапазоны, Unicode, пустые и очень длинные значения, повторяющиеся ключи.
  • Временные тесты: позднее событие, out-of-order, справочник до и после event_time, backdated correction.
  • Идемпотентность: одно событие отправляется несколько раз и создаёт один эффективный результат.
  • Fault injection: timeout, 429, 5xx, cache outage, malformed response, падение после API call и до commit, восстановление из checkpoint.
  • Нагрузочные тесты: пик, hot key, высокая кардинальность, cache stampede, рост DLQ, медленный downstream.
  • Data-quality тесты: completeness, join miss, конфликт источников, distribution drift и stale reference.
  • LLM eval: фиксированный набор, evidence coverage, доля abstain, нарушения схемы, hallucination и согласие reviewer.
  • Privacy/security тесты: redaction, field allowlist, отсутствие секретов и PII в логах/DLQ, удаление во всех слоях.

Как перейти со старого конвейера

Безопасная миграция не начинается с перезаписи основной таблицы.

  1. Описать потребителей, источник истины, SLO и текущую семантику полей.
  2. Сохранить raw и добавить side table или новый topic с версией результата.
  3. Запустить shadow computation: считать новую версию, но не выдавать её пользователям.
  4. Сравнить completeness, расхождения, latency, стоимость и долю review на репрезентативной истории.
  5. Включить canary для части tenants, entity IDs или трафика. Dual-read логирует отличия.
  6. Запустить backfill отдельными partition и quota с меньшим приоритетом, чем live-трафик.
  7. Переключить current pointer после quality gates; сохранить быстрый возврат.
  8. Сверить counts, hashes и downstream read-back. Старый путь отключать после окна наблюдения.

Backfill должен иметь собственные rate limit и cost budget. Миллиард старых записей, выпущенных в ту же очередь, легко исчерпают API quota и увеличат latency живого трафика.

Что мониторить после запуска

Метрики должны показывать не только доступность сервиса, но и смысл результата.

Область Метрики и сигналы
Поток input/output rate, queue lag по event time, oldest pending age, throughput по partition
Задержка P50/P95/P99 по enricher и end-to-end, timeout share, time in queue
Надёжность success/error, retries, DLQ size и drain rate, circuit state, replay failures
Качество completeness, null/unknown, join miss, invalid enum/range, duplicate/conflict rate
Свежесть now - source_as_of, CDC lag, age of materialized result, stale fallback share
Кэш hit/miss, origin QPS, evictions, hot keys, stampede и negative-cache share
Стоимость calls/event, cost/event, batch size, LLM tokens/event, дневной budget и quota remaining
Модели model/prompt/schema version, distribution drift, abstain, hallucination sample, human override
Governance lineage coverage, PII policy failures, retention/deletion backlog, unowned contracts

Alert должен вести к действию. Например, рост join miss при нормальной инфраструктуре указывает на изменение ключа или задержку reference topic. Рост cache hit вместе с ростом stale age говорит о слишком длинном TTL или сломанной инвалидации. Нулевая DLQ при падении success rate может означать потерю записей, а не идеальную обработку.

Как выбрать архитектуру

Начните с пяти решений:

  1. Какое бизнес-действие использует поле? От этого зависят обязательность и цена ошибки.
  2. К какому времени должен относиться факт? Текущий, на event_time или известный системе тогда.
  3. Какой режим укладывается в SLO? Request-time, async stream, batch или гибрид.
  4. Что произойдёт при отказе? Stop, pending, partial, stale или manual review.
  5. Как доказать и отменить результат? Provenance, version, evidence, raw, replay и rollback.

Для большинства развивающихся продуктов разумная отправная точка — асинхронное write-time обогащение: raw в PostgreSQL или durable log, отдельные workers, versioned side table, идемпотентный upsert, ограниченные retry, DLQ, Redis cache и явные pending/partial/ready. Kafka/Flink стоит добавлять, когда подтверждены требования к потоку, replay, нескольким потребителям и stateful temporal joins, а не ради названия архитектуры.

Вывод

Enrichment-архитектура превращает «добавим ещё одно поле» в управляемый производственный процесс. Её качество определяют не число коннекторов и моделей, а ясные временные семантики, сохранённый raw, идемпотентность, происхождение каждого атрибута, контролируемая деградация, версионирование и возможность backfill.

Главный выбор проходит между задержкой, свежестью, стоимостью и согласованностью. Request-time даёт свежесть и сильную связанность, batch — экономию и задержку, streaming local join — производительность и сложное состояние, async write-time — масштабирование и eventual consistency. Правильная схема делает этот компромисс видимым и измеримым для потребителей.

Автор статьи

Елена Кравцова — Редактор и автор статей
Елена Кравцова

Редактор и автор статей

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

Вопросы и ответы

Нет. Для умеренной нагрузки достаточно PostgreSQL, таблицы заданий, workers и Redis. Kafka оправдан, когда нужны длительный replayable log, несколько независимых потребителей, partitioned ordering и большой поток. Выбирайте транспорт по требованиям, а не по термину.

ETL описывает движение и порядок обработки: extract, transform, load. Enrichment — конкретная цель одной или нескольких transform-стадий: добавить полезный контекст. Один ETL-конвейер может очищать, нормализовать, агрегировать и обогащать данные.

В RAG найденные фрагменты обогащают контекст запроса к модели на request time, но исходная бизнес-запись обычно не меняется. Подготовка RAG-индекса — извлечение сущностей, тегов, языка, прав доступа и embeddings — может быть полноценным enrichment-конвейером.

Да, пока это допускают цель, стоимость и правила хранения. Raw нужен для аудита, исправления ошибки и replay новой версии. Для персональных данных задайте срок, удаление и распространение запроса на все копии, включая кэш и DLQ.

Нет. Exactly-once внутри stream processor или транзакционного sink не всегда охватывает сторонний side effect. Падение после успешного API call и до checkpoint создаёт повтор. Используйте idempotency key поставщика, dedup store и сверку результата.

TTL должен быть не длиннее допустимого окна устаревания и учитывать частоту изменения источника. Проверьте hit rate, origin QPS и stale age на реальном трафике. Для важных справочников сочетайте TTL с инвалидацией по CDC и защитой от cache stampede.

Можно, если контракт перечисляет готовые и ожидаемые поля, а потребитель умеет работать со статусом partial. Для обязательного risk-score, юридического ограничения или другого blocking-атрибута quality gate должен задержать действие, а не маскировать отсутствие значением null.

Зарегистрировать новую версию, проверить совместимость, считать её в shadow mode и переключать потребителей постепенно. Breaking change лучше выпускать как отдельную major-версию или topic/table. Старый результат сохраняют на окно rollback.

Самооценка модели не является откалиброванной вероятностью. Надёжнее требовать evidence spans, проверять результат детерминированными правилами и измерять точность на размеченном наборе. Для рискованных случаев вводят abstain и human review.

Используйте отдельную очередь или consumer group, меньший приоритет, отдельные concurrency и cost budgets, ограничение по partitions и паузы при росте live lag. Сначала прогоните маленький диапазон, сверив counts, ошибки и стоимость.

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

Смотрите также

Поделиться

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

Оставьте комментарий

Войдите, чтобы присоединиться к обсуждению