АрхитектураМикросервисы и интеграцииАрхитектор программных систем

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

Новый потребитель подключается спустя месяц к потоку событий. Какой тип событийного контракта позволит ему восстановить актуальное состояние?

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

  1. Достаточно ли хранить только последнее событие по каждой сущности?

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

Такое хранение подходит для контракта текущего состояния, если это явно заявлено и потребителю не требуется replay истории. Для воспроизводимого контракта нужна полная история за гарантированный период либо эквивалентный механизм снимков с доступными событиями после снимка.

  1. Можно ли считать повторное проигрывание событий безопасным без идемпотентности потребителя?

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

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

  1. Что произойдёт, если порядок событий нарушен при восстановлении состояния?

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

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