Microservice: FoodTrackerAsyncService, обработка и публикация сообщений в Kafka

Архитектурная спецификация: Ingestion API Producer (Write Path)

Published

June 11, 2026

WarningОграничение публичной документации

В открытом доступе представлена демонстрационная версия метода. В настоящей публичной документации отображены не все шаги, технические сценарии и приватные эндпоинты для системы цифровых симуляторов бизнес-процессов.

  • Полная спецификация метода: Будет доступна только во внутреннем контуре разработки (Confluence / Swagger Enterprise).

Создадим обработчик промежуточного шага, внутренный метод/ы сервиса, который будет читать из кафки сообщения, выполнять бизнес-логику и отправлять в другой топик кафки сообщение результат обработки.

1 Общее назначение и архитектурная роль компонента

Компонент foodTrackerAsyncService функционирует как горизонтально масштабируемый gRPC-продюсер. Его единственная задача — принять запрос от внешнего клиента (мобильного приложения или симулятора), преобразовать его в событие и отправить в шину данных Kafka со строгим партиционированием по ключу group_id.

2 Схема архитектуры компонента (Producer Flow)

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

3 Системная спецификация: Ingestion API Producer (Write Path)

3.1 Общее назначение и архитектурная роль компонента

Компонент foodTrackerAsyncService выполняет роль высокопроизводительного, горизонтально масштабируемого входного шлюза (Command Producer) для регистрации фактов жизненного цикла продуктов. Компонент функционирует в качестве gRPC-сервера на языке Go. Ключевые архитектурные решения:

  • Асинхронный неблокирующий ввод: Входные транспортные хэндлеры полностью изолированы от транзакционной логики базы данных 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”}}