Messaging и routing
Правила Kafka topics, message keys, envelopes и Event Router.
Messaging и routing
Kafka и Event Router являются внутренней частью платформы. Их контракты должны быть управляемыми, но не должны протекать в consumer-facing API.
Topic contract
Каждый topic регистрируется в Topic Catalog и содержит:
- имя и назначение;
- domain, source и target;
- producer и consumers;
- owner;
- message schema/envelope;
- message key и partitioning rationale;
- retention;
- retry и DLT/DLQ topics;
- data classification.
Topic не создаётся ad hoc
Перед созданием проверяется naming policy и отсутствие коллизии. Topic без owner, schema и retention не готов к использованию.
Naming
Имя формируется из утверждённых смысловых сегментов: domain, direction/purpose, system и operation/event. Версия добавляется только согласно Kafka naming policy, а не произвольно каждой командой.
Специальные retry, reply, orchestration и DLT topics должны следовать отдельным каноническим шаблонам.
Message key
Key выбирается по требуемому порядку обработки и распределению нагрузки. Он должен:
- быть стабилен для одной логической сущности;
- обеспечивать порядок там, где он действительно нужен;
- не создавать hot partition;
- извлекаться по документированному правилу;
- не содержать чувствительные данные открытым текстом.
Envelope
Минимальный envelope переносит:
{
"consumerId": "platform-resolved-id",
"externalRequestId": "request-attempt-id",
"traceparent": "00-...",
"idempotencyKey": "business-request-key",
"occurredAt": "2026-07-12T14:30:00+03:00",
"payload": {}
}Конкретная schema утверждается отдельно; пример показывает разделение идентификаторов.
Routing rule
Rule связывает входной API context/method и классификацию flow с целевым topic. Rule имеет owner, priority, условия, target и историю изменений.
Router не должен содержать скрытую доменную оркестрацию. Если решение требует нескольких бизнес-шагов, маршрут передаёт управление Kestra.
Event Router baseline на 2026-08-19
Последний Confluence export уточняет текущую реализацию Router:
- источником ingress-событий является WSO2;
- входной REST-вызов имеет семантику fire-and-forget, а публикация в Kafka выполняется синхронно;
- target topic выбирается по request parameters и JSON Logic rule;
- перед добавлением route должны быть зарегистрированы входной и выходной topics в Topic Catalog;
callbackUrlпереключает обработку в async flow; для него результат отправляется черезcallback.out;- правила хранятся в БД и кешируются Caffeine на 24 часа; принудительное обновление выполняется через
POST /reload.
Sync/async semantics требуют уточнения
Экспорт одновременно описывает sync-чтение reply topic и утверждает, что без callbackUrl чтение из Kafka не требуется. До исправления первичного документа конкретный response flow следует подтверждать по route и end-to-end тесту.
Outbox и восстановление публикации
При недоступности Kafka Router сохраняет зашифрованный payload в event_outbox. Retry scheduler запускается каждые 5 секунд и использует retry_count, max_retries, next_retry_at, last_attempt_at и статусы NOT_PROCESSING, RETRYING, DELETE.
Payload в outbox защищается AES-256-GCM. После чтения он расшифровывается перед Kafka publish. Записи со статусом DELETE, которым больше трёх дней, очищаются через pgcron.
Kafka key и partitions
В текущем описании key и partitioning помечены как неиспользуемые: при kafka_key = null сообщение направляется в partition 0, хотя topic описан с тремя partitions. Это временное ограничение, а не рекомендуемый целевой дизайн. Перед включением распределения нагрузки необходимо определить key, проверить ordering guarantees и обновить Router Rules List вместе с Topic Catalog.
Delivery semantics
Проектируйте обработчики с предположением at-least-once delivery. Даже если отдельный компонент обещает stronger semantics, adapter и Consumer callback должны безопасно обрабатывать повтор.