5.0 KiB
Интерфейсы сервиса 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.