graph TD
%% Определение стилей компонентов
classDef transport fill:#E3F2FD,stroke:#0D47A1,stroke-width:1px;
classDef asyncSvc fill:#FFF9C4,stroke:#FBC02D,stroke-width:1px;
classDef kafka fill:#D9EAD3,stroke:#274E13,stroke-width:2px;
subgraph "Транспортный слой (Входной контур)"
gRPC[gRPC API Handlers]:::transport
end
subgraph "Бизнес-логика (Async Core)"
AsyncSvc[foodTrackerAsyncService]:::asyncSvc
KP[KafkaPublisher Interface]:::asyncSvc
end
subgraph "Шина данных (Выходной контур)"
T1["topic: purchase-events"]:::kafka
T2["topic: cooking-events"]:::kafka
T3["topic: consumption-events"]:::kafka
T4["topic: waste-events"]:::kafka
end
%% Потоки данных
gRPC -->|1. Вызов команд| AsyncSvc
AsyncSvc -->|2. Очистка токенов и маршалинг в JSON| KP
KP -->|3. Key: group_id| T1
KP -->|3. Key: group_id| T2
KP -->|3. Key: group_id| T3
KP -->|3. Key: group_id| T4
Microservice: FoodTrackerAsyncService, обработка и публикация сообщений в Kafka
Архитектурная спецификация: Ingestion API Producer (Write Path)
В открытом доступе представлена демонстрационная версия метода. В настоящей публичной документации отображены не все шаги, технические сценарии и приватные эндпоинты для системы цифровых симуляторов бизнес-процессов.
- Полная спецификация метода: Будет доступна только во внутреннем контуре разработки (Confluence / Swagger Enterprise).
1 Общее назначение и архитектурная роль компонента
Компонент foodTrackerAsyncService функционирует как горизонтально масштабируемый gRPC-продюсер. Его единственная задача — принять запрос от внешнего клиента (мобильного приложения или симулятора), преобразовать его в событие и отправить в шину данных Kafka со строгим партиционированием по ключу group_id.
2 Схема архитектуры компонента (Producer Flow)
3 Системная спецификация: Ingestion API Producer (Write Path)
- Асинхронный неблокирующий ввод: Входные транспортные хэндлеры полностью изолированы от транзакционной логики базы данных PostgreSQL. Они выполняют быструю валидацию контрактов, очищают метаданные от чувствительной информации (dto.AccessToken = ““), преобразуют структуры домена в JSON-строки и атомарно сбрасывают их в брокер Apache Kafka. Это гарантирует минимальный и стабильный p99 времени ответа (duration_ms) gRPC-клиентам и симуляторам даже при пиковых нагрузках массового прогона.
- Партиционирование по бизнес-контексту: Публикация всех событий в шину данных выполняется со строгой фиксацией ключа Key = group_id. Брокер Kafka гарантирует, что сообщения с идентичным ключом всегда распределяются в одну и ту же партицию. Это обеспечивает строгую хронологическую последовательность обработки команд консьюмерами в рамках одной домашней группы (исключая состояние гонки, когда команда списания или готовки обгоняет команду закупки).
- Сквозная трассировка (Distributed Tracing): Транспортный слой оснащен перехватчиком, который извлекает маркер сквозного логирования x-request-id из gRPC-метаданных от симулятора/клиента. При отсутствии заголовка генерируется новый UUID v4. Этот маркер принудительно пробрасывается в контексте во все структурированные логи slog и в метаданные сообщений Kafka для сквозного аудита.
4 Спецификация Входных gRPC-интерфейсов (Транспортный контракт)
Транспортный слой обрабатывает данные в рамках спецификации FoodTrackerServiceAsync. В случае инфраструктурных сбоев (например, сетевой разрыв с кластером брокера) или ошибок сериализации, внутренние исключения Go преобразуются в канонические коды ошибок пакета google.golang.org/grpc/codes.
4.1 Метод TrackPurchase (Регистрация закупки по чеку)
Принимает состав фискального чека, очищает авторизационный токен и пушит payload в топик purchase-events.
Входные параметры (pb.PurchaseRequest): Содержит уникальный номер чека receipt_no, итоговую стоимость total_price, строковый UUID group_id и массив позиций items. Если receipt_no или group_id пустые, сервер возвращает статус codes.InvalidArgument.
Выходные параметры (pb.EmptyResponse): При успешной передаче в брокер возвращается пустой ответ со статусом codes.OK.
4.2 Метод TrackCooking (Регистрация приготовления блюда)
Фиксирует создание нового составного блюда из набора существующих ингредиентов.
- Входные параметры (pb.CookingRequest): Передает UUID группы group_id, ID создаваемого блюда dish_item_id, название dish_name, выходные метрики yield_units/yield_weight_grams и массив используемого сырья ingredients (source_item_id, used_units, used_weight_grams).
4.3 Методы TrackConsumption и TrackWaste (Списание остатков)
Принимают плоские идентификаторы для утилизации или употребления порций продукта в пищу.
- При пустом значении item_id или отрицательных значениях объемов (units / weight) отдается статус codes.InvalidArgument.
5 Сценарий взаимодействия при публикации команд (Sequence Diagram)
Ниже представлена диаграмма обработки gRPC-вызова, работы интерцептора логирования, маршалинга структуры в JSON и отправки сообщения в Kafka кластер.
%%{init: {
'theme': 'base',
'themeVariables': {
'actorBkg': '#E3F2FD',
'actorBorder': '#546E7A',
'actorTextColor': '#0D47A1',
'rectBkg': '#FFF9C4',
'rectBorder': '#FBC02D',
'noteBkgColor': '#F3E5F5',
'noteBorderColor': '#7E57C2',
'noteTextColor': '#311B92',
'signalColor': '#2C3E50',
'signalLineColor': '#2C3E50'
}
}}%%
sequenceDiagram
autonumber
actor Client as Внешний Клиент / Симулятор
participant Interceptor as IngestionLoggingInterceptor
participant Server as foodTrackerAsyncService
participant KP as KafkaPublisher (Client)
Client->>Interceptor: Вызов метода gRPC (например, TrackPurchase) с MD [x-request-id]
activate Interceptor
note over Interceptor: Шаг 2: Извлечение x-request-id из метаданных.<br/>При отсутствии — генерация uuid.NewV4()
note over Interceptor: Шаг 3: Инициализация JSON-логгера slog с фиксацией request_id и метода
Interceptor->>Server: Перенаправление контекста и DTO во внутренний хэндлер
activate Server
note over Server: Шаг 5: Безопасность — dto.AccessToken = ""<br/>Маршалинг структуры DTO в бинарный JSON payload
Server->>KP: Вызов PublishEvent(ctx, "purchase-events", groupID, payload)
activate KP
KP-->>Server: Успешный возврат nil (Подтверждение ACK от Kafka зафиксировано)
deactivate KP
Server-->>Interceptor: Возврат пустой структуры ответа и статуса ошибки (nil)
deactivate Server
note over Interceptor: Шаг 9: Расчет duration_ms и запись лога JSON Lines в os.Stdout
Interceptor-->>Client: Ответ EmptyResponse + gRPC Status Code (codes.OK)
deactivate Interceptor
5.1 Таблица расшифровки шагов сценария асинхронного логирования команд
| Шаг | Действие | Параметры / Запросы / DTO | Код ошибки (canonical_code) |
HTTP / gRPC статус |
|---|---|---|---|---|
| 1 (Client -> Interceptor) | Симулятор или клиент отправляет DTO команды, передавая маркер трассировки в заголовках. | gRPC Metadata Context:x-request-id: “uuid-v4-trace-id” | INGEST_INVALID_ARGUMENT | codes.InvalidArgument (400) |
| 2-3 (Interceptor -> Interceptor) | Внутренняя логика: Регистрация request_id, инициализация потока логирования и сброс лога о старте операции в stdout. | slog JSON Line Payload:{ “msg”: “Async request started”, “request_id”: “uuid”, “method”: “TrackPurchase” } | Нет | codes.OK (200) |
| 4 (Interceptor -> Server) | Интерцептор пробрасывает модифицированный контекст на уровень асинхронного Go-сервиса. | Внутренний вызов Go:s.TrackPurchase(ctx, userID, groupID, dto) | Нет | codes.OK (200) |
| 5 (Server -> Server) | Внутренняя логика: Стирание токенов из структуры во избежание утечки в логи брокера. Сериализация DTO. | Внутренняя функция:json.Marshal(dto) | JSON_MARSHAL_FAILURE | codes.Internal (500) |
| 6 (Server -> KP) | Сервис вызывает метод клиента Kafka, передавая имя целевого топика и ключ партиционирования. | Аргументы метода:PublishEvent(ctx, “purchase-events”, groupID, payload) | KAFKA_PUBLISH_TIMEOUT | codes.Internal (500) |
| 7 (KP -> Server) | Драйвер брокера подтверждает фиксацию сообщения в коммит-логе партиции и возвращает управление. | Kafka Wire Network ACK | KAFKA_CLUSTER_DOWN | codes.Internal (500) |
| 8-9 (Server -> Interceptor) | Проброс успешного статуса выполнения, замер общего времени удержания сетевого gRPC-соединения. | slog JSON Line Payload:{ “level”: “INFO”, “duration_ms”: 1.8, “request_id”: “trace-id” } | Нет | codes.OK (200) |
| 10 (Interceptor -> Client) | Клиент получает ответ об успешной регистрации команды. | Задача ушла на обработку в очередь. |
6 Справочник Обсервабилити (Структура логов Producer-контура для ClickHouse)
Каждая строка лога, сбрасываемая в os.Stdout, содержит строго типизированные поля для сборщика Vector. При аварийных ситуациях в ClickHouse фиксируются следующие метрики:
| 1. Код (canonical_code) | 2. Уровень (log_level) | 3. Статусы (transport_statuses) | 4. Получатель (error_target) | 5. JSON для Фронтенда (ui_payload) | 6. Метрики для ClickHouse (observability_json) |
|---|---|---|---|---|---|
| INGEST_INVALID_ARGUMENT | WARN | {“grpc”: 3, “http”: 400} | EXTERNAL_SYSTEM | {“error”: “Неверный формат данных чека или группы”} | {“metric”: “bad_request”, “labels”: {“source”: “simulator”}} |
| JSON_MARSHAL_FAILURE | ERROR | {“grpc”: 13, “http”: 500} | INTERNAL_SYSTEM | null | {“metric”: “serialization_err”, “labels”: {“service”: “ingestion_async”}} |
| KAFKA_PUBLISH_TIMEOUT | ERROR | {“grpc”: 13, “kafka_ack”: 0} | INTERNAL_SYSTEM | {“error”: “Шина данных временно перегружена”} | {“metric”: “kafka_timeout”, “labels”: {“partition_key”: “group_id”}} |
| KAFKA_CLUSTER_DOWN | CRITICAL | {“grpc”: 13, “kafka_ack”: -1} | INFRASTRUCTURE | {“error”: “Сервис временно недоступен”} | {“metric”: “kafka_broker_dead”, “labels”: {“cluster”: “production”}} |