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







