camunda-to-amqpЗамена AMQP-прокладки на Kafka domain events для orders
Цель: убрать сервис-прокладку camunda-to-amqp из production. user-task-creator становится единой точкой интеграции с Camunda (External Tasks + Kafka domain events, без AMQP для orders).
backoffice → AMQP → camunda-to-amqp → Camunda REST
Latency overhead и 4 точки отказа.
Отдельный сервис, helm-чарт, on-call ответственность.
DLQ через rascal есть, но без retry/backoff в коде worker-ов.
RabbitMQ только для orders + Kafka для остальных доменов (Invoice / Shipment / Pi).
user-task-creator — единственное место, работающее с Camunda (External Tasks + Kafka events, без прокладок).
5 узлов, каждая связь — через AMQP-hop
camunda-to-amqp: TaskCreation (External Task), ShipmentFinished (External Task → REST к BER), orders-create (AMQP → startOrderProcess), orders-update (AMQP → cancel + terminate).
3 узла + Kafka, без промежуточного брокера для orders
domain events (order_created / order_updated)
producer в BER (уже работает для PI/Shipment)
consumer в UTC (ServerKafka + @EventPattern)
Создание заказа → запуск Camunda-процесса
Идемпотентность: hasProcessInstance pre-check перед startOrderProcess. Camunda /start не гарантирует уникальность businessKey — поэтому явный pre-check обязателен. 400/409 (race) → soft-ack.
Отмена заказа → отмена тасок + терминирование процессов
camunda-to-amqp реальную работу делал только status === 'Cancel'. Completed — no-op (getTasksByOrderId вызывался, результат отбрасывался). В этом PR воспроизводим AS IS. Нужно ли терминировать процессы при Completed — открытый вопрос Q1 для аналитиков.
completeProcessInstance = DELETE /process-instance/{id} (hard terminate, не BPMN-completion).
EventsService.sendDomainEvent({key, event}) — уже работает в проде для PI/Shipment. key = order.number (сквозной, для partitioning). EventBusClient (uuid v4) — deprecated, миграция out of scope.
Confluent .KafkaJS-драйвер. Throw → seek-back → redelivery. Новая зависимость для user-task-creator (T0). Single-handler + instanceof dispatch (паттерн из integration-service).
Только id, number, status — для AS IS-поведения. Остальные поля добавляем по мере необходимости (не засоряем контракт).
Как EntityChangedEventPayload (консистентно с Invoice/Shipment/Pi). Примечание: before может быть подвержен race condition.
hasProcessInstance перед startOrderProcess — защита от дублей, replay, zombie processes. Существующий getProcessInstance throw-ит → пишем новую обёртку с мягким 404.
Handler переезжает из camunda-to-amqp. External Task (не Message Event). Меняется delivery-семантика: silent at-most-once → at-least-once (throw → incident). PUT /shipments/{id} идемпотентен ✅.
| # | Риск | Митигация |
|---|---|---|
| R1 | Нет outbox — при Kafka-down заказ зависнет без Camunda-процесса | Retry 3× (exp backoff) в EventsService + метрика orders_publish_errors_total + reconciliation job (parking lot). Принимаем как риск. |
| R2 | businessKey не гарантирует дедупликацию | Pre-check hasProcessInstance перед start. 400/409 race → soft-ack. |
| R3 | Новый failure mode: Kafka-зависимость | Non-commit offset + alert на consumer lag (promaas-shared KFKR101). Business-ошибки → soft-ack. |
| R4 | Poison message стопает партицию (нет auto-DLQ) | Manual offset move через kafka-consumer-groups.sh + runbook. Alert на lag-growth. |
| R5 | Смена семантики ShipmentFinished (at-most-once → at-least-once) | Осознанное улучшение. PUT /shipments/{id} идемпотентен ✅. |
| R6 | Completed — no-op (возможный латентный баг) | Воспроизводим AS IS. Вынесено в Q1 для аналитиков. |
Каждый шаг обратим (feature flags). Между релизами — 1-2 недели стабилизации.
ries-modules: OrderEntity + events. BER: EventsService (flag=OFF). UTC: Kafka consumer (flag=OFF) + ShipmentFinished handler. Эффект: 0.
rollback: revert PR
BER: dual-write (flag=ON). UTC: consumer (flag=ON). camunda-to-amqp: Orders workers OFF. Наблюдаем 1-2 недели.
rollback: flag OFF
BER: AMQP-publish OFF. Удалить camunda-to-amqp/ целиком. Удалить broker config.
irreversible
AMQP payload не меняется на этапе dual-write — старый формат { entity, operationType, entityBody } сохраняется до Релиза 3. Любое изменение сломает оставшийся в строю camunda-to-amqp.
getTasksByOrderId, cancelTasksByOrderId) + @idp/nestjs-kafka dep 0.5dИтого: ~9.5d разработки, ~10-12d календарных (2 дева в параллель).
/metrics:9464): orders_*_handled_total, orders_*_errors_total{reason}, orders_start_skipped_totalpromaas-shared P0174.yaml (правило KFKR101, метрика kafka_consumergroup_lag). Добавить наш consumer group в T4.KafkaRetriableException в event-mode = обычный throw. Retry-лимита и auto-DLQ нет → poison message стопает партицию до ручного offset move (runbook).
AS IS — no-op. Воспроизводим AS IS. Решить с бизнесом → follow-up PR если "да".
DELETE /process-instance/{id} — terminate, не BPMN-completion. Сохраняем AS IS.
Итог: план готов к реализации после разрешения Q1 (Completed-status). Три фатальные технические ошибки из аудита исправлены (fromBeginning в consumer-config, контрактный пакет расширяется, @idp/nestjs-kafka устанавливается). Cutover в 3 релиза — каждый шаг обратим до Релиза 3.
📋 Полную документацию см. в [RIES] Убить сервис camunda-to-amqp 🔪.md