Паттерн Saga через черги повідомлень для мікросервісів
Реалізація паттерну 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 днів |
Зв'яжіться з нами для консультації — ми проаналізуємо вашу архітектуру та запропонуємо оптимальне рішення. Замовте впровадження, і ми допоможемо уникнути типових помилок. Отримайте надійну розподілену транзакцію без блокувань.







