iac/apps/documentations/dps-message-hub.ENDPOINTS.md

5.0 KiB
Raw Blame History

Интерфейсы сервиса dps-message-hub

Версия: 0.1.0

Документ описывает интерфейсную поверхность сервиса: потребляемые топики Kafka и обращения к внешним зависимостям (PostgreSQL).

dps-message-hub — это чистый Kafka-воркер на FastStream. У него нет HTTP/REST API (в коде нет FastAPI/Flask/aiohttp, health-роутов и т.п.), поэтому файла openapi.yaml для сервиса нет. Исходящих HTTP-запросов к другим сервисам он тоже не выполняет — единственный получатель данных — база PostgreSQL.

Как устроено взаимодействие

Единое FastStream-приложение (dps_message_hub.infra.app:get_app, src/dps_message_hub/infra/app.py) объединяет:

  • Kafka — потребитель сообщений (KafkaBroker + KafkaRouter), топик assets (src/dps_message_hub/features/assets/interface/kafka.py);
  • PostgreSQL — пул соединений psycopg (AsyncConnectionPool), открывается/закрывается в lifespan (src/dps_message_hub/infra/lifespan.py, infra/database.py).

Каждое сообщение проходит через middleware Retry (src/dps_message_hub/interface/middleware.py): при исключении обработка повторяется бесконечно с экспоненциальной задержкой (1 << min(10, retry_count) секунд). Коммит оффсета ручной — у потребителя auto_commit=False.

Kafka-потребители (входящие сообщения)

Реальное имя топика задаётся переменной DPS_MESSAGE_HUB_KAFKA__TOPICS (маппинг логического имени assets в имя топика Kafka). Формат сообщения — BrokerMessageDto (src/dps_message_hub/infra/dto.py): поля schema_version, model, sender, type, body, diff, timestamp, trace_id (alias xtraceId), user_id, tenants, tags.

Топик (логич.) group_id Offset reset auto_commit Обработчик Назначение
assets dps_assets_consumer earliest false assets Обновление событий разметки (markup_event) по изменению ассета

Диспетчеризация внутри обработчика assets по полю type (src/dps_message_hub/features/assets/interface/kafka.py); body разбирается в модель AssetUpdate (поля id, attributes[]), diff — произвольный dict:

type сообщения Условие (diff) Действие Назначение
model_updated diff is None Пропуск (нет изменений)
model_updated в diff есть ключ resource_id unlink_asset(asset_id) Отвязка ассета: деактивация его событий разметки
model_updated в diff есть ключ attributes update_markup_events(asset_id, attributes) Пересоздание событий разметки по обновлённым атрибутам
model_deleted unlink_asset(asset_id) Отвязка ассета при удалении модели

Бизнес-логика (src/dps_message_hub/features/assets/application/asset.py): для обновляемых атрибутов существующие активные markup_event деактивируются (inactivation_time), затем в транзакции вставляются новые записи. Тип атрибута (поле type) определяет целевую колонку значения (value_int/value_float/value_string/value_option_id/unit_option_id/value_dt/value_date) — маппинг в features/assets/domain/entities.py.

Внешние зависимости (инфраструктура)

Зависимость Назначение
Kafka Источник сообщений (топик assets). Подключение SASLScram512, при заданном KAFKA__SSL_CAFILE — по SASL_SSL
PostgreSQL Хранилище данных (psycopg async, таблица markup_event). Операции: SELECT, UPDATE, массовая вставка COPY ... FROM STDIN (features/assets/infra/database.py)

Исходящие HTTP-запросы к внешним сервисам

Отсутствуют. Сервис не содержит HTTP-клиентов (httpx/aiohttp/requests) и не обращается к другим сервисам по HTTP; все побочные эффекты — запись в PostgreSQL.