Как правильно организовать связь между сервисами смешанного типа и оповещать пользователя о завершении процесса?
Я вовлечен в создание системы состоящей из нескольких клиентов и сервисов. В качестве клиентов через которые пользователь взаимодействует с системой выступают плагины для настольных приложений и WebUI. Весь список сервисов я указывать не буду, перечислю только те, что важны для обозначения ситуации:
Kong API Gatewayза которым всё прячется;- Сервис работы с базой данных (
MongoDB), написан наFastAPI, представляет из себя обычныйREST API; - Сервис обработки информации, сейчас это тоже
FastAPI, порождающийCeleryзадачи, в которых я запрашиваю файлы из сервиса БД и обрабатываю их; Celeryсервис, который принимает задачи и выполняет их (в качестве брокера используетсяRedis);- Сервис написанный на
NestJSобслуживающийWebUI.
В текущей версии стандартный порядок работы следующий:
- Пользователь с помощью плагина отправляет
HTTPзапрос с данными к сервису БД; - Сервис БД валидирует информацию, если всё хорошо, пишет в БД;
- После создания новой сущности в
MongoDB, из сервиса БД отправляется запрос к сервису обработки информации (для конвертации в формат по умолчанию). Ответ не ожидается; - Сервис обработки информации получает запрос и если никто до этого
уже не проводил эту же операцию, создает
Celeryзадачу; Celeryсервис производит конвертацию, помещает результат в репозиторий и на этом всё;- Пользователь открывает
WebUIи видит сохранённую в БД сущность (сервисNestJSделает запрос к сервису БД), в окне просмотра отображается модель полученная в процессе конвертации (сервисNestJSделает запрос к сервису обработки информации); - Если пользователь хочет конвертировать сущность не в формат по умолчанию, а в какой-нибудь специфический, то он делает запрос к эндпоинту сервиса обработки информации;
- Сервис обработки информации создаёт
Celeryзадачу и отдаёт в ответе пользователю еёtask_id; - Сервис
NestJSпериодически спрашивает сервис обработки информации о состоянииCeleryзадачи; - Когда конвертация заканчивается и статус задачи меняется на
SUCCEEDEDпользователь запрашивает результат у сервиса обработки данных.
Было принято решение перейти от подобной системы, использующей для общения между сервисами HTTP запросы к общению через очередь сообщений. В качестве инструмента для реализации был выбран NATS Streaming. Сейчас, на этапе продумывания изменений возникают вопросы, на которые я, в силу своей неопытности, не могу объективно ответить самому себе.
Например:
Если я после создания сущности в БД, помещу сообщение о событии в общую шину NATS, то могу с лёгкостью в новом сервисе произвести нужные действия и так же по их окончанию поместить новое сообщение оповещающее об окончании процесса обработки, тут всё хорошо, НО… Что, если кто-то захочет конвертировать сущность в другой формат? Что тогда? Допустим я оставлю REST сервис обработки данных для инициации этой процедуры, тогда в нём я должен поместить в NATS сообщение о событии “была запрошена конвертация в такой-то формат” вместо создания задачи Celery? Сервис “слушающий” подобные “события” начнёт конвертацию и… И… и что? Пользователь-то по HTTP запрос сделал, эндпоинт создал “эвент” и ответил, что мол ждите, а чего ждать-то? Как оповестить клиент, что конвертация закончилась или, что произошла ошибка? В текущем варианте с Celery есть task_id, по которому можно проверять текущее положение дел, а с новой системой как? Возникает мысль, что комбинировать REST сервисы с сервисами общающимися через MQ не очень хорошая идея. Но что тогда делать? Может создать отдельный сервис который держит постоянную связь с клиентом через WebSocket соединение, такой же как NestJS для WebUI только для других клиентов? Пользователь инициировал конвертацию, сервис слушает подобные события и при ответе оповещает пользователя по WS? Что тогда брать для реализации подобных вещей? Может вообще пускать всех через тот самый NestJS сервис который обслуживает WebUI, а чего, пусть и с плагинами держит соединение?! Но следом возникает вопрос авторизации, как в NATS фильтровать только события, нужные текущему пользователю? Если N пользователей приконнектилось, каждый запустил конвертацию своей сущности из БД в нужный формат, создалось N событий в NATS, вида entity.converted, то что, все сервисы слушающие шину событий на предмет завершения конвертации примут ВСЕ эвенты и вручную надо тело каждого сообщения сериализовать и смотреть моя ли это сущность была конвертирована?
В итоге, главный вопрос такой, что можно почитать, чтобы в голове всё встало на свои места?