Назначение#
Цель создания#
Platform V Synapse Event Mesh (#EM) — это комплекс интеграционных компонентов для потоковой обработки, передачи и мониторинга событий в цифровых системах, обеспечивающий гибкую работу с событиями между современными сервисами и ИИ-агентами.
Позволяет решить следующие задачи:
Synapse Event Monitoring system (Mayak) (EDMN)#
Программный компонент Synapse Event Monitoring system (Mayak) EDMN представляет собой инструмент для сбора метрик с технических сервисов и отображения собранной информации в виде информативных дашбордов в UI, что позволяет осуществлять мониторинг состояния доменов и своевременно реагировать на возникающие проблемы.
Компонент предназначен для:
сбора метрик с СПО и ППО событийного сегмента;
отслеживания состояния метрик с помощью конфигурируемых триггеров;
отображения собранных метрик и их состояния на конфигурируемом дашборде.
Сервис управления событийными доменами (EDMS)#
Позволяет клиенту заказывать интеграцию при помощи пользовательского интерфейса, реализует жизненный цикл заявок и позволяет развертывать интеграции на транспортном слое.
Integration Portal (EMIP)#
Позволяет клиенту:
обеспечивать в рамках GenAI-трансформации создание коммуникационной среды для мультиагентных систем с развязкой по производительности и недоступности;
давать возможность AI-агентам работать с событиями, порождаемыми источниками;
давать возможность агентам искать и подключаться к каналам, содержащим необходимую информацию, на естественном языке;
предоставлять инструменты для автоматического создания асинхронных интеграционных цепочек.
Synapse Cloud Event Streaming Processing (EVPC)#
Компонент представляет собой тот же самый универсальный потоковый обработчик, адаптированный для работы в среде Kubernetes и других платформ, основанных на нем.
Установка возможна только в контейнеризированной среде.
Pod обработчика функционально аналогичен Job Apache Flink и настраивается с помощью конфигурации, которая передается при запуске.
Помимо DSL, также можно выполнять расширенное конфигурирование с использованием языка JavaScript, что предоставляет более подготовленным пользователям расширенные возможности настройки логики обработки.

Обработчик поддерживает нативное подключение к транспорту на основе Apache Kafka.
Для подключения к другим типам транспорта, таким как RabbitMQ или Apache ActiveMQ Artemis, используется адаптер EVTA.
Кроме того, EVTA может выступать в роли egress-шлюза, что особенно полезно в инсталляциях, где требуется использование Istio.
Сервис обратного отсчета потоковой обработки (EVPT)#
Компонент представляет собой вспомогательный инструмент для работы с таймерами, который принимает REST-запросы для установки таймеров и выполняет указанный REST-вызов по истечении заданного времени.
Установка сервиса возможна как на виртуальных машинах, так и в контейнерах. Сервис поддерживает два режима работы: персистентный, с хранением заданий в базе данных на основе PostgreSQL, и неперсистентный, с хранением заданий в оперативной памяти.
В рамках потоковой обработки EVPT используется для отсчета окон агрегации в облачном потоковом обработчике EVPC, где он разворачивается в Kubernetes. При необходимости EVPT может быть использован вне контекста потоковой обработки в качестве самостоятельного планировщика задач.
Reactive stream adapter (EVTA)#
Программный компонент Reactive stream adapter (EVTA) представляет собой адаптер интеграционного слоя, предназначенный для:
обеспечения обмена данных между топиками Platform V Corax / Apache Kafka и очередями сообщений (MQ) с использованием REST и gRPC без трансформации данных;
репликации данных между топиками Platform V Corax / Apache Kafka и очередями сообщений (MQ) в любой комбинации;
валидации передаваемых сообщений по заданной JSON/AVRO/XSD схеме.
Подключение к транспорту по REST/gRPC#
Некоторые системы, являющиеся потенциальными источниками и потребителями событий при переходе к событийной архитектуре, не заточены под взаимодействие непосредственно с транспортом по нативному протоколу и не могут перестроить свою архитектуру в силу каких-либо причин. Чаще всего такие системы умеют взаимодействовать по REST или gRPC, и EVTA в данном случае представляет собой «прослойку» между подобными системами и транспортным слоем.

Межтранспортная репликация#
Адаптер EVTA можно использовать для репликации данных между топиками Platform V Corax / Apache Kafka и очередями сообщений (MQ) в любой комбинации:

Важно понимать, что передаваться будут данные конкретных настроенных в адаптере потоков (топиков или очередей), а не дублироваться весь кластер целиком.
Внутренняя функциональность EVTA#
Адаптер имеет набор внутренних функций, который может быть доработан и расширен. Например, при передаче сообщений EVTA уже позволяет:
выполнять валидацию сообщений по заданным XSD/ AVRO/ JSON-схемам;
выполнять преобразование тела сообщения;
выполнять фильтрацию заголовков;
позволяет гибко настраивать логирование прохождения сообщения на всех своих внутренних этапах.
Варианты развертывания#
EVTA может инсталлироваться как на ВМ, так и в контейнеризированной среде. В дистрибутиве присутствуют DevOps-скрипты, обеспечивающие автоматизацию развертывания.
Сервис передачи событий (EVTD)#
Программный компонент «Сервис передачи событий» (EVTD) предназначен для:
обеспечения дополнительных сервисных функций при отправке и получении сообщений через Platform V Corax / Apache Kafka;
обеспечения функций администрирования Apache Kafka.
Функциональность EVTD#
Компонент EVTD представляет собой набор различных перехватчиков (интерсепторов) – это вспомогательные модули, позволяющие выполнить над событием какую-либо дополнительную логику. Kafka-клиенты имеют встроенный механизм подключения перехватчиков, поэтому перехватчики можно подключить к продюсеру и консьюмеру в их конфигурации соответственно.
Перехватчики на продюсер отрабатывают после получения события из прикладного кода, но до отправки в транспортный слой. На консьюмер наоборот – после получения события из транспортного слоя, но до передачи в прикладной код. Прикладной разработчик может создавать собственные перехватчики с абсолютно любой функциональностью. Можно также подключить цепочку перехватчиков, которые будут выполняться последовательно.

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

В EVTD реализованы:
перехватчики валидации (XML, JSON, AVRO, отправка событий аудита по результатам валидации);
перехватчики подписи (х509 и ОТТ), также поставляются в виде де-/сериализаторов;
Отдельные специфические виды перехватчиков:
dsl-перехватчик – позволяет преобразовывать тело сообщений при помощи языка DSL;
latency-перехватчик – измеряет задержки прохождения события;
json-consumer-interceptor – выполняет фильтрацию JSON-сообщений, используя jsonPath;
json-extractor-interceptor – может достать поле из сообщения в JSON-формате и добавить его в заголовок.
Варианты развертывания#
Компонент EVTD устанавливается на ВМ.
Сервис потоковой обработки событий (EVTP)#
Компонент представляет собой кластер Apache Flink, управляемый с помощью Zookeeper, и предназначен для установки исключительно на виртуальные машины (ВМ).
Apache Flink состоит из двух ключевых компонентов: JobManager и TaskManager. Кластер может включать несколько JobManager, из которых активным является только один, а остальные выполняют роль резервных для обеспечения отказоустойчивости, а также несколько TaskManager, которые работают одновременно, повышая общую производительность системы.
Каждый обработчик, или «задание потоковой обработки» (Job), делится на шаги (Task). Job управляется JobManager, а Task(s) выполняются на TaskManager(s).
Apache Flink поддерживает концепцию «параллелизма»: если запустить Job с уровнем параллелизма больше 1, обработка будет выполняться в несколько потоков, что значительно увеличивает производительность.

При использовании Apache Flink в чистом виде или при реализации логики обработки внутри систем-участников требуется отдельно кодировать каждое преобразование, привлекать ресурсы разработки, проводить пересборку системы, тестирование и внедрение новой версии ПО.
Наше решение устраняет эти сложности за счет универсального потокового обработчика, который представляет собой Job Apache Flink, способный принимать логику работы в виде конфигурации, распознавать ее и выполнять на указанных потоках данных.
Логика обработки задается через конфигурационные файлы, что исключает необходимость написания кода. Для конфигурирования мы разработали упрощенный скриптовый язык (DSL), доступный даже пользователям с минимальной алгоритмической подготовкой.
Конфигурации оформляются в виде файлов, и после их создания все, что нужно сделать для развертывания обработки - это запустить универсальный обработчик с этой конфигурацией.
Состав продукта#
В состав Platform V Synapse Event Mesh (#EM) входят следующие компоненты:
Synapse Event Monitoring system (Mayak) (EDMN)#
Компонент, предназначен для сбора метрик с технических сервисов и отображения собранной информации в виде информативных дашбордов в пользовательском интерфейсе, что позволяет осуществлять мониторинг состояния доменов и своевременно реагировать на возникающие проблемы.
Сервис управления событийными доменами (EDMS)#
Компонент, предназначен для создания заявок на публикацию новых событий, подписку на существующие события и репликацию существующего события в другой домен через пользовательский интерфейс.
Integration Portal (EMIP)#
Компонент, решает интеграционную задачу работы AI-агентов с событиями и между собой в мультиагентной среде, а так же помогает классическим автоматизированным системам облегчить интеграцию с событийным слоем.
Synapse Cloud Event Streaming Processing (EVPC)#
Компонент, предназначенный для потоковой обработки событий, поступающих из одного или нескольких событийных доменов. Результатом потоковой обработки событий является новый поток событий, поставщиком которого является сервис потоковой обработки событий. Потоковый обработчик на базе Apache Flink. Компонент предназначен для развертывания только на ВМ.
Сервис обратного отсчета потоковой обработки (EVPT)#
Компонент, предназначенный для предоставления возможности работы с расписанием задач. Задача представляет собой запланированный единичный или повторяющийся HTTP/HTTPs-запрос. Вспомогательный сервис, может использоваться как для отслеживания окон агрегации облачного потокового обработчика, так и самостоятельно для прикладных нужд. Компонент предназначен для развертывания на ВМ и в облачной среде исполнения Kubernetes или OpenShift.
Reactive stream adapter (EVTA)#
Компонент, предназначенный для обмена данных между топиками Platform V Corax / Apache Kafka, брокерами сообщений ArtemisMQ и IBM MQ, возможность работы с транспортом (Platform V Corax / Apache Kafka, ArtemisMQ, IBM MQ) с использованием REST и gRPC взаимодействий, валидация сообщений по заданным XSD/ AVRO/ JSON-схемам
Сервис передачи событий (EVTD)#
Компонент, предназначенный для обеспечения дополнительных сервисных функций при отправке и получении сообщений через Platform V Corax / Apache Kafka; обеспечение функций администрирования Apache Kafka
Сервис потоковой обработки событий (EVTP)#
Компонент, предназначенный для потоковой обработки событий, поступающих из одного или нескольких событийных доменов. Результатом потоковой обработки событий является новый поток событий, поставщиком которого является сервис потоковой обработки событий. Потоковый обработчик на базе Apache Flink. Компонент предназначен для развертывания только на ВМ.