1 / 12
RIES · Low-priority backlog

🔪 Убить сервис camunda-to-amqp

Замена AMQP-прокладки на Kafka domain events для orders

Цель: убрать сервис-прокладку camunda-to-amqp из production. user-task-creator становится единой точкой интеграции с Camunda (External Tasks + Kafka domain events, без AMQP для orders).

−1
сервис в проде
−1
helm-чарт
−1
on-call команда
3
релиза (безопасный cutover)
2 / 12

Содержание

01Зачем убивать сервис
02AS-IS архитектура
03TO-BE архитектура
04Поток данных: Order created
05Поток данных: Order updated (Cancel)
06Ключевые решения (ADR)
07Риски и митигации
08Cutover — 3 релиза
09Задачи и оценка
10Открытые вопросы
11Observability
12Итог
3 / 12

Зачем убивать сервис

🔗

4 звена в цепочке orders

backoffice → AMQP → camunda-to-amqp → Camunda REST

Latency overhead и 4 точки отказа.

⚠️

Single point of failure

Отдельный сервис, helm-чарт, on-call ответственность.

DLQ через rascal есть, но без retry/backoff в коде worker-ов.

🔀

Дублирование инфраструктуры

RabbitMQ только для orders + Kafka для остальных доменов (Invoice / Shipment / Pi).

🎯

Единая точка интеграции

user-task-creator — единственное место, работающее с Camunda (External Tasks + Kafka events, без прокладок).

4 / 12

AS-IS архитектура

5 узлов, каждая связь — через AMQP-hop

Camunda
TaskCreationTopic · ShipmentFinished · REST API
⬇ ⬇
camunda-to-amqp ❌
3 worker-а: proxy + orders create + orders update
RabbitMQ
4 очереди для orders + tasks
⬇ ⬇ ⬇
user-task-creator    backoffice-entity-repository
4 потока через camunda-to-amqp: TaskCreation (External Task), ShipmentFinished (External Task → REST к BER), orders-create (AMQP → startOrderProcess), orders-update (AMQP → cancel + terminate).
5 / 12

TO-BE архитектура

3 узла + Kafka, без промежуточного брокера для orders

Camunda
External Tasks (pull) · REST API (startOrderProcess, DELETE)
⬇ External Tasks    ⬆ HTTP
user-task-creator
createTask · processShipmentMessage · OrdersKafkaController
⬆ consume (Kafka)    ⬇ HTTP (CamundaClient · EntityRepoService)
Kafka
ries-system-experience-domain-events-v1
key = order.number
⬆ publish
backoffice-entity-repository
EventsService.sendDomainEvent({ key, event })
   camunda-to-amqp
Kafka

domain events (order_created / order_updated)

EventsService

producer в BER (уже работает для PI/Shipment)

@idp/nestjs-kafka

consumer в UTC (ServerKafka + @EventPattern)

6 / 12

Поток: Order created

Создание заказа → запуск Camunda-процесса

backoffice
POST /orders
DB write
Kafka
order_created
key=order.number
user-task-creator
OrdersKafkaController
hasProcessInstance?
pre-check
404 (нет процесса) → startOrderProcess (POST /start, businessKey=number)     200 (есть процесс) → skip + warn (защита от дублей)

Идемпотентность: hasProcessInstance pre-check перед startOrderProcess. Camunda /start не гарантирует уникальность businessKey — поэтому явный pre-check обязателен. 400/409 (race) → soft-ack.

7 / 12

Поток: Order updated (Cancel)

Отмена заказа → отмена тасок + терминирование процессов

backoffice
PUT /orders/:id/rms
status=Cancel
Kafka
order_updated
{before, after}
user-task-creator
status === 'Cancel' → getTasksByOrderId → cancelTasksByOrderId → loop DELETE /process-instance/{id}
status === 'Completed'no-op (AS IS — см. Q1)
др. status → soft-ack
⚠ AS IS-особенность: в camunda-to-amqp реальную работу делал только status === 'Cancel'. Completed — no-op (getTasksByOrderId вызывался, результат отбрасывался). В этом PR воспроизводим AS IS. Нужно ли терминировать процессы при Completedоткрытый вопрос Q1 для аналитиков.

completeProcessInstance = DELETE /process-instance/{id} (hard terminate, не BPMN-completion).
8 / 12

Ключевые архитектурные решения

ADR #7

EventsService вместо EventBusClient

EventsService.sendDomainEvent({key, event}) — уже работает в проде для PI/Shipment. key = order.number (сквозной, для partitioning). EventBusClient (uuid v4) — deprecated, миграция out of scope.

ADR #1

@idp/nestjs-kafka для consumer

Confluent .KafkaJS-драйвер. Throw → seek-back → redelivery. Новая зависимость для user-task-creator (T0). Single-handler + instanceof dispatch (паттерн из integration-service).

ADR #3

OrderEntity — минимум полей

Только id, number, status — для AS IS-поведения. Остальные поля добавляем по мере необходимости (не засоряем контракт).

ADR #8

Payload {before, after}

Как EntityChangedEventPayload (консистентно с Invoice/Shipment/Pi). Примечание: before может быть подвержен race condition.

ADR #9

Pre-check guard

hasProcessInstance перед startOrderProcess — защита от дублей, replay, zombie processes. Существующий getProcessInstance throw-ит → пишем новую обёртку с мягким 404.

ADR #10

ShipmentFinished → user-task-creator

Handler переезжает из camunda-to-amqp. External Task (не Message Event). Меняется delivery-семантика: silent at-most-once → at-least-once (throw → incident). PUT /shipments/{id} идемпотентен ✅.

9 / 12

Риски и митигации

#РискМитигация
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 для аналитиков.
10 / 12

Cutover — 3 релиза

Каждый шаг обратим (feature flags). Между релизами — 1-2 недели стабилизации.

1
День 0

Подготовка

ries-modules: OrderEntity + events. BER: EventsService (flag=OFF). UTC: Kafka consumer (flag=OFF) + ShipmentFinished handler. Эффект: 0.

rollback: revert PR

2
День 7-14

Активация

BER: dual-write (flag=ON). UTC: consumer (flag=ON). camunda-to-amqp: Orders workers OFF. Наблюдаем 1-2 недели.

rollback: flag OFF

3
День 21-28

Cleanup

BER: AMQP-publish OFF. Удалить camunda-to-amqp/ целиком. Удалить broker config.

irreversible

AMQP payload не меняется на этапе dual-write — старый формат { entity, operationType, entityBody } сохраняется до Релиза 3. Любое изменение сломает оставшийся в строю camunda-to-amqp.

11 / 12

Задачи и Observability

🛠 Задачи (Релиз 1)

  • T0 — контракт (EntityRepoService: getTasksByOrderId, cancelTasksByOrderId) + @idp/nestjs-kafka dep 0.5d
  • T1 — ries-modules: OrderEntity + events 1d
  • T2 — BER: EventsService publish 1d
  • T3 — UTC: consumer + handlers + CamundaClient 2d
  • T4 — helm + tests + E2E + promaas-shared KFKR101 1.5d

Итого: ~9.5d разработки, ~10-12d календарных (2 дева в параллель).

📊 Observability

  • Counters (telemetry-стек, /metrics:9464): orders_*_handled_total, orders_*_errors_total{reason}, orders_start_skipped_total
  • Consumer lagне из библиотеки. Источник: Strimzi/kafka-exporter → promaas-shared P0174.yaml (правило KFKR101, метрика kafka_consumergroup_lag). Добавить наш consumer group в T4.
  • Datadog alerts на рост error-rate + lag-growth
Error handling (confluent-драйвер): throw → seek-back → redelivery (non-commit offset). KafkaRetriableException в event-mode = обычный throw. Retry-лимита и auto-DLQ нет → poison message стопает партицию до ручного offset move (runbook).
12 / 12

Открытые вопросы и итог

🙋 Открытые вопросы

Q1 (аналитикам): Completed — терминировать процесс?

AS IS — no-op. Воспроизводим AS IS. Решить с бизнесом → follow-up PR если "да".

Q2 (known): completeProcessInstance = terminate

DELETE /process-instance/{id} — terminate, не BPMN-completion. Сохраняем AS IS.

✅ Подтверждено при челлендже

  • POST /tasks/cancel идемпотентен (state='new' filter)
  • PUT /shipments/{id} идемпотентен (хронология — отдельный эндпоинт)
  • AMQP consumer audit: 0 потребителей incoming-message-queue
  • order.number — unique (Mongo + Postgres)
  • Статусы: только Underway / Completed / Cancel

Итог: план готов к реализации после разрешения Q1 (Completed-status). Три фатальные технические ошибки из аудита исправлены (fromBeginning в consumer-config, контрактный пакет расширяется, @idp/nestjs-kafka устанавливается). Cutover в 3 релиза — каждый шаг обратим до Релиза 3.

📋 Полную документацию см. в [RIES] Убить сервис camunda-to-amqp 🔪.md