Метод (Event Handler) ProcessParsedReceiptEvent

Домен: AID | Сервис: purchase-ollama-worker | Тип: Kafka Consumer

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

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

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

Общее описание

Асинхронный метод-обработчик (Event Handler) ProcessParsedReceiptEvent развернут внутри микросервиса purchase-ollama-worker. Он отвечает за семантический разбор (NER/Парсинг) очищенной текстовой строки и извлечение из нее структурированных сущностей (название товара, объем/вес, процент жирности/алкоголя, бренд, стоимость).

Метод полностью изолирован от синхронного HTTP-цикла мобильного приложения, не имеет REST-эндпоинта и активируется исключительно при поступлении нового сообщения в брокер. Результатом работы является публикация унифицированного паспорта расходов в веерный топик дистрибуции для сервисов инвентаря и биллинга.

Метод решает следующие задачи:

  1. Динамическая адаптация промпта (Prompt Localization): Подмена системного контекста и правил парсинга ИИ в зависимости от метаданных языка (app_lang), переданных шлюзом.
  2. ИИ-распознавание сущностей (NER Structured Outputs): Взаимодействие с LLM-контейнером (Llama) для гарантированного получения строго типизированного ответа по Pydantic-схеме (JSON-массив).
  3. Генерация паспорта расходов (Fan-out Preparation): Пуш структурированного шаблона в топик дистрибуции bpds.mdm.out.product.templated для параллельной вычитки сервисами fridge-service и billing-service.

Протокол взаимодействия и триггеры (Event Contract)

  • Имя метода: ProcessParsedReceiptEvent
  • Тип метода: Event Handler / Kafka Consumer & Producer
  • Брокер сообщений: Apache Kafka
  • Входной топик (Subscription): bpds.inventory.out.receipt.parsed
  • Выходной топик (Publication): bpds.mdm.out.product.templated
  • Формат данных: application/json

Спецификация входного события (Ingress Message Payload)

Воркер ожидает на входе очищенный текст, успешно прошедший аудит безопасности:

Поле Тип Обязательный Описание Пример значения
trace_id String Да Сквозной ID запроса для трассировки логов req-manual-text-99aa
user_id String Да Уникальный идентификатор пользователя usr_9876
home_group_id String Да ID домашней группы для привязки инвентаря group_abc123
clean_text_input String Да Очищенная от мата и спама текстовая строка Молоко Юрта в ауле 3.2% 1.5 л 1029 тенге
app_lang String Да Двухбуквенный код языка для локализации промпта ru

Пример JSON-события из топика (Payload):

{
  "trace_id": "req-manual-text-99aa",
  "user_id": "usr_9876",
  "home_group_id": "group_abc123",
  "clean_text_input": "Молоко Юрта в ауле 3.2% 1.5 л 1029 тенге",
  "app_lang": "ru"
}

Диаграмма последовательности (Mermaid)

sequenceDiagram
    autonumber
    
    participant K as Брокер: Apache Kafka Cluster
    participant Worker as purchase-ollama-worker (AID)
    participant Llama as Llama Container: NER Parsing (AID)

    %% ЭТАП 1: ИИ-РАСПОЗНАВАНИЕ И СИНТЕЗ JSON
    K->>Worker: Шаг 1: Вычитка: bpds.inventory.out.receipt.parsed (Текст + app_lang)
    activate Worker
    note over Worker: Шаг 2: Адаптация системного промпта<br/>под локализацию app_lang
    
    critical Шаг 3: API-запрос к LLM для извлечения сущностей
        Worker->>Llama: API: Запрос Structured Outputs (Pydantic-схема)
        activate Llama
        Llama-->>Worker: Шаг 4: Возврат структурированного JSON-массива продуктов
        deactivate Llama
    option Ошибка инференса нейросети (Таймаут / Сбой VRAM)
        Note over Worker: Экстренный прерванный цикл:<br/>Событие маркируется как сбойный лог
    end
    
    note over Worker: Шаг 5: Валидация паспорта расходов в памяти
    
    Worker->>K: Шаг 6: Пуш в топик: bpds.mdm.out.product.templated
    deactivate Worker


Расшифровка шагов

Шаг Действие Параметры / Запросы / DTO Ошибки (Исключения / Статусы)
1 (K -> Worker) Воркер, находясь в бесконечном цикле поллинга брокера, вычитывает новое событие с очищенным текстом из топика цензуры. Kafka Consumer Poll Request
Topic: bpds.inventory.out.receipt.parsed
Payload:
{ "trace_id": "req-manual-text-99aa", "clean_text_input": "Молоко Юрта 3.2%...", "app_lang": "ru" }
Бизнес-ошибки отсутствуют
2 (Worker -> Worker) Внутренняя логика: Микросервис извлекает из памяти или кэша конфигурационный файл системного промпта, соответствующий языку app_lang (правила выделения валют, объемов, стоп-символов). Внутренний метод сборщика контекста:
PromptEngine.compile(payload.app_lang)
IAD-NER-400 (передан неподдерживаемый app_lang, движок не смог подобрать схему локализации промпта)
3 (Worker -> Llama) Воркер отправляет структурированный gRPC/HTTP-запрос к контейнеру Ollama/Llama. Запрос содержит жесткую JSON-схему (Structured Outputs), обязывающую ИИ вернуть строго типизированный ответ. HTTP POST /api/chat (Ollama Options)
Payload (NERRequestDTO):
{ "model": "llama3:ner-fridge", "messages": [{"role": "system", "content": "..."}, {"role": "user", "content": "Молоко Юрта..."}], "format": "json" }
Бизнес-ошибки отсутствуют
4 (Llama -> Worker) Нейросеть завершает инференс и возвращает массив распознанных продуктов, разбитый по атрибутам, готовый к десериализации в бэкенд DTO. Response Body (NERResponseDTO):
{ "items": [{ "product_name": "Молоко Юрта", "fat_percentage": 3.2, "volume_liters": 1.5, "price": 1029.0, "currency": "KZT" }] }
IAD-NER-422 (ИИ вернул невалидную структуру, галлюцинировал или не смог распарсить сырую строку в массив товаров)
5 (Worker -> Worker) Внутренняя валидация: Воркер проверяет полученный массив в памяти, рассчитывает общую сумму чека (паспорт расходов) и проверяет обязательные поля для инвентаризации. Внутренняя проверка:
Pydantic.validate_values(items)
IAD-VAL-400 (ИИ вернул пустой массив продуктов или отрицательные цены/объемы)
6 (Worker -> K) Воркер обогащает JSON метаданными пользователя и отправляет готовый шаблон продуктов в веерный топик Kafka для распределения по сервисам инвентаря и биллинга. Kafka Message (Topic: bpds.mdm.out.product.templated):
Key: "usr_9876"
Payload:
{ "trace_id": "req-manual-text-99aa", "user_id": "usr_9876", "home_group_id": "group_abc123", "parsed_items": [...] }
Бизнес-ошибки отсутствуют