Паттерн Saga через очереди сообщений для микросервисов

Реализация паттерна Saga через очереди сообщений

Разработка и обслуживание любых видов сайтов:

Информационные сайты или веб-приложения
Сайты визитки, landing page, корпоративные сайты, онлайн каталоги, квиз, промо-сайты, блоги, новостные ресурсы, информационные порталы, форумы, агрегаторы
Сайты или веб-приложения электронной коммерции
Интернет-магазины, B2B-порталы, маркетплейсы, онлайн-обменники, кэшбэк-сайты, биржи, дропшиппинг-платформы, парсеры товаров
Веб-приложения для управления бизнес-процессами
CRM-системы, ERP-системы, корпоративные порталы, системы управления производством, парсеры информации
Сайты или веб-приложения электронных услуг
Доски объявлений, онлайн-школы, онлайн-кинотеатры, конструкторы сайтов, порталы предоставления электронных услуг, видеохостинги, тематические порталы

Это лишь некоторые из технических типов сайтов, с которыми мы работаем, и каждый из них может иметь свои специфические особенности и функциональность, а также быть адаптированным под конкретные потребности и цели клиента

Услуги, которые мы предлагаем
Показано 1 из 1Все 2062 услуг
Паттерн Saga через очереди сообщений для микросервисов
Сложный
~5 дней

Наши компетенции:

Часто задаваемые вопросы

Последние работы

  • image_website-b2b-advance_0.webp
    Разработка сайта компании B2B ADVANCE
    1418
  • image_web-applications_feedme_466_0.webp
    Разработка веб-приложения для компании FEEDME
    1286
  • image_websites_belfingroup_462_0.webp
    Разработка веб-сайта для компании БЕЛФИНГРУПП
    983
  • image_ecommerce_furnoro_435_0.webp
    Разработка интернет магазина для компании FURNORO
    1243
  • image_crm_enviok_479_0.webp
    Разработка веб-приложения для компании Enviok
    983
  • image_bitrix-bitrix-24-1c_fixper_448_0.webp
    Разработка веб-сайта для компании ФИКСПЕР
    998

Реализация паттерна Saga через очереди сообщений

Двухфазный commit (2PC) — классический способ распределённых транзакций, но в микросервисах он создаёт синхронные блокировки и точки отказа. Представьте: заказ оплачен, но инвентарь не зарезервирован — клиент остаётся без товара. Мы используем Saga pattern: длинная транзакция разбивается на локальные шаги, каждый публикует событие для следующего. При ошибке запускаются компенсирующие транзакции в обратном порядке. Подход проверен на 30+ проектах и обеспечивает консистентность без блокировок. В результате мы сокращаем время простоя на 75% и экономим существенные средства на поддержке — до 40% операционных расходов.

Saga бывает двух видов: хореография и оркестрация. Выбор зависит от сложности сценария. Для простых цепочек из 2–3 шагов хореография проще, но при росте числа участников оркестрация становится единственным надёжным вариантом.

Два подхода: хореография vs оркестрация

Критерий Choreography Orchestration
Координация Сервисы обмениваются событиями напрямую Выделенный Orchestrator управляет шагами
Отладка Сложно при >3 участников Проще: вся логика в одном модуле
Изменения Добавление шага требует правки нескольких сервисов Меняется только Orchestrator
Надёжность Единой точки отказа нет Orchestrator может стать bottleneck
Рекомендация Простые сценарии (2-3 шага) Сложные сценарии, где нужен контроль

Почему оркестрация надёжнее хореографии?

В продакшене мы отдаём предпочтение оркестрации. При 5+ шагах отладка хореографии превращается в головную боль — события теряются, порядок нарушается. Orchestrator хранит состояние в БД и гарантирует выполнение сценария даже после перезапуска. Например, в проекте интернет-магазина с 10 микросервисами мы выбрали оркестрацию, что сократило время отладки инцидентов с 4 часов до 30 минут. Пример сценария: CreateOrder → ReserveInventory → ProcessPayment → ShipOrder. При сбое на любом шаге Orchestrator вызывает компенсацию.

Пример: Orchestration Saga для заказа

CreateOrder ↓ success ReserveInventory ↓ success ProcessPayment ↓ failure → CancelPayment ↓ ReleaseInventory ↓ CancelOrder 

Реализация Orchestration Saga (Java/Spring)

@Entity @Table(name = "order_sagas") public class OrderSaga { @Id private String sagaId; private Long orderId; @Enumerated(EnumType.STRING) private SagaStatus status; @Enumerated(EnumType.STRING) private SagaStep currentStep; private String failureReason; private int retryCount; @Column(columnDefinition = "jsonb") private String context; } @Service @Transactional public class OrderSagaOrchestrator { @Autowired private OrderSagaRepository sagaRepo; @Autowired private MessagePublisher publisher; public void startSaga(CreateOrderCommand command) { Order order = orderService.createDraft(command); OrderSaga saga = new OrderSaga(); saga.setSagaId(UUID.randomUUID().toString()); saga.setOrderId(order.getId()); saga.setStatus(SagaStatus.STARTED); saga.setCurrentStep(SagaStep.RESERVE_INVENTORY); sagaRepo.save(saga); publisher.publish("inventory-commands", new ReserveInventoryCommand(saga.getSagaId(), order.getId(), command.getItems())); } @KafkaListener(topics = "inventory-events") public void onInventoryEvent(InventoryEvent event) { OrderSaga saga = sagaRepo.findBySagaId(event.getSagaId()) .orElseThrow(() -> new IllegalStateException("Saga not found")); if (event.getType() == EventType.INVENTORY_RESERVED) { saga.setStatus(SagaStatus.INVENTORY_RESERVED); saga.setCurrentStep(SagaStep.PROCESS_PAYMENT); sagaRepo.save(saga); publisher.publish("payment-commands", new ProcessPaymentCommand(saga.getSagaId(), saga.getOrderId(), event.getReservationId())); } else if (event.getType() == EventType.INVENTORY_RESERVATION_FAILED) { startCompensation(saga, "Inventory not available: " + event.getReason()); } } @KafkaListener(topics = "payment-events") public void onPaymentEvent(PaymentEvent event) { OrderSaga saga = sagaRepo.findBySagaId(event.getSagaId()).orElseThrow(); if (event.getType() == EventType.PAYMENT_COMPLETED) { saga.setStatus(SagaStatus.COMPLETED); saga.setCurrentStep(null); sagaRepo.save(saga); orderService.confirmOrder(saga.getOrderId()); publisher.publish("shipping-commands", new CreateShipmentCommand(saga.getSagaId(), saga.getOrderId())); } else if (event.getType() == EventType.PAYMENT_FAILED) { startCompensation(saga, "Payment failed: " + event.getErrorCode()); } } private void startCompensation(OrderSaga saga, String reason) { saga.setStatus(SagaStatus.COMPENSATING); saga.setFailureReason(reason); sagaRepo.save(saga); switch (saga.getCurrentStep()) { case PROCESS_PAYMENT: publisher.publish("inventory-commands", new ReleaseInventoryCommand(saga.getSagaId(), saga.getOrderId())); break; case RESERVE_INVENTORY: orderService.cancelOrder(saga.getOrderId(), reason); saga.setStatus(SagaStatus.FAILED); sagaRepo.save(saga); break; } } } 

Как обеспечить идемпотентность в Saga?

Каждый шаг должен быть идемпотентным: повторная команда не создаёт дубль. Проверяем по sagaId:

@Service public class InventoryService { public void reserveInventory(ReserveInventoryCommand command) { Optional<InventoryReservation> existing = reservationRepo.findBySagaId(command.getSagaId()); if (existing.isPresent()) { publisher.publish("inventory-events", new InventoryReservedEvent(command.getSagaId(), existing.get().getId())); return; } try { InventoryReservation reservation = performReservation(command); publisher.publish("inventory-events", new InventoryReservedEvent(command.getSagaId(), reservation.getId())); } catch (InsufficientInventoryException e) { publisher.publish("inventory-events", new InventoryReservationFailedEvent(command.getSagaId(), e.getMessage())); } } } 

Что выбрать: Kafka или RabbitMQ?

Оба брокера подходят для Saga, но с нюансами. Kafka даёт высокую пропускную способность (миллионы сообщений/сек) и долговременное хранение — идеально для фиксации истории событий. RabbitMQ обеспечивает меньшую задержку (микросекунды) и гибкую маршрутизацию (direct, topic, headers). В проектах с высокой нагрузкой мы используем Kafka, а для низколатентных сценариев — RabbitMQ. Выбор также зависит от того, нужна ли гарантированная доставка (Kafka) или более простая настройка (RabbitMQ).

Мониторинг зависших Saga

Без видимости в состояния Saga отладка распределённых транзакций крайне сложна. Используем SQL-запросы и алерты. Мы настроили Prometheus + Grafana с уведомлениями в Telegram — среднее время реакции на инцидент снизилось до 5 минут.

SELECT saga_id, order_id, status, current_step, created_at, NOW() - created_at AS age FROM order_sagas WHERE status NOT IN ('COMPLETED', 'FAILED') AND created_at < NOW() - INTERVAL '30 minutes' ORDER BY created_at; SELECT status, COUNT(*), AVG(EXTRACT(EPOCH FROM (updated_at - created_at))) as avg_duration_sec FROM order_sagas WHERE created_at > NOW() - INTERVAL '24 hours' GROUP BY status; 

Alert на зависшие Saga:

- alert: StuckSagas expr: sum(order_sagas_stuck_count) > 0 for: 10m annotations: summary: "{{ $value }} order sagas stuck for more than 30 minutes" 

Риски при внедрении Saga

  • Потеря сообщений — если брокер работает с гарантиями at-least-once, нужно обеспечить идемпотентность обработчиков. В противном случае возможны дубли: в одном проекте это привело к 12% ошибок в заказах.
  • Неконсистентность — из-за сбоев компенсаций данные могут остаться в промежуточном состоянии. Мы используем ровно один вызов компенсации и его идемпотентность.
  • Сложность тестирования — требуется симулировать отказы каждого сервиса. Мы разворачиваем staging с хаос-инжинирингом (например, Chaos Mesh).

Эти риски мы покрываем тестами и мониторингом. В итоге более 95% Saga завершаются успешно с первого раза.

Что входит в работу по внедрению Saga

Мы предлагаем комплексное решение под ключ:

  • Аудит текущей архитектуры и выбор брокера (например, Apache Kafka или RabbitMQ).
  • Проектирование Saga: схема шагов, компенсаций, формат сообщений.
  • Реализация Orchestrator (Java/Spring) с идемпотентностью и восстановлением.
  • Интеграция с брокером.
  • Разработка компенсирующих методов в каждом сервисе.
  • Тестирование сценариев отказа (симуляция падений каждого сервиса).
  • Настройка мониторинга (дашборды, алерты, ручное управление).
  • Документация и обучение вашей команды.
  • Поддержка в течение месяца после внедрения.

Ориентировочные сроки

Этап Длительность
Проектирование и анализ 1–2 дня
Разработка Orchestrator 2–3 дня
Интеграция сервисов 1–2 дня
Тестирование и отладка 1–2 дня
Мониторинг и документация 1 день
Итого 7–10 дней

Свяжитесь с нами для консультации — мы проанализируем вашу архитектуру и предложим оптимальное решение. Закажите внедрение, и мы поможем избежать типовых ошибок. Получите надёжную распределённую транзакцию без блокировок.