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







