Saga Pattern with Message Queues for Microservices

Implementing the Saga Pattern via Message Queues

Development and maintenance of all types of websites:

Informational websites or web applications
Business card websites, landing pages, corporate websites, online catalogs, quizzes, promo websites, blogs, news resources, informational portals, forums, aggregators
E-commerce websites or web applications
Online stores, B2B portals, marketplaces, online exchanges, cashback websites, exchanges, dropshipping platforms, product parsers
Business process management web applications
CRM systems, ERP systems, corporate portals, production management systems, information parsers
Electronic service websites or web applications
Classified ads platforms, online schools, online cinemas, website builders, portals for electronic services, video hosting platforms, thematic portals

These are just some of the technical types of websites we work with, and each of them can have its own specific features and functionality, as well as be customized to meet the specific needs and goals of the client.

Our competencies:

Frequently Asked Questions

Latest works

  • image_web-applications_feedme_466_0.webp
    Development of a web application for FEEDME
    1287
  • image_ecommerce_furnoro_435_0.webp
    Development of an online store for the company FURNORO
    1244
  • image_crm_enviok_479_0.webp
    Development of a web application for Enviok
    983
  • image_crm_chasseurs_493_0.webp
    CRM development for Chasseurs
    1034
  • image_website-sbh_0.webp
    Website development for SBH Partners
    1108
  • image_website-_0.webp
    Website development for Red Pear
    555

Implementing the Saga Pattern via Message Queues

Two-phase commit (2PC) is a classic approach for distributed transactions, but in microservices it creates synchronous locks and single points of failure. Imagine: an order is paid, but inventory is not reserved — the customer ends up without the product. The Saga pattern is essential for managing distributed transactions in microservices using message queues. We use the Saga pattern: a long-running transaction is split into local steps, each publishing an event for the next. On error, compensating transactions run in reverse order. This approach has been proven on 30+ projects, providing consistency without locking. As a result, we reduce downtime by 75% and save significant operational costs — up to 40% (our clients typically save $100,000–$150,000 annually on support costs). Our implementation service costs $15,000 on average.

We cover both Saga orchestration and choreography, using message queues like Kafka and RabbitMQ, with emphasis on idempotency and monitoring.

Saga comes in two flavors: choreography and orchestration. The choice depends on scenario complexity. For simple chains of 2–3 steps, choreography is simpler, but as the number of participants grows, orchestration becomes the only reliable option. In production, our orchestrator handles up to 5,000 concurrent sagas with 99.99% consistency rate.

Choreography or Orchestration: Which Approach to Choose?

Comparison of Choreography and Orchestration
Criteria Choreography Orchestration
Coordination Services exchange events directly A dedicated Orchestrator manages steps
Debugging Complex with >3 participants Easier: all logic in one module
Changes Adding a step requires modifying multiple services Only the Orchestrator changes
Reliability No single point of failure Orchestrator can become a bottleneck
Recommendation Simple scenarios (2-3 steps) Complex scenarios where control is needed

Why Orchestration Is More Reliable Than Choreography

In production, we favor orchestration. With 5+ steps, debugging choreography becomes a headache — events get lost, order breaks. Orchestrator debugging is 8x faster than choreography. The Orchestrator stores state in a database and guarantees scenario execution even after restart. For example, in an e-commerce project with 10 microservices, we chose orchestration, reducing incident debugging time from 4 hours to 30 minutes (8x improvement). Example scenario: CreateOrder → ReserveInventory → ProcessPayment → ShipOrder. On failure at any step, the Orchestrator triggers compensation.

Example: Orchestration Saga for Orders

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

Implementation of 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; } } } 

Ensuring Idempotency in Saga

Each step must be idempotent: a repeated command does not create a duplicate. We check by 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())); } } } 

Choosing Between Kafka and RabbitMQ

Both brokers suit Saga, but with nuances. Kafka offers high throughput (millions of messages/sec) and long-term storage — ideal for event sourcing. RabbitMQ provides lower latency (microseconds) and flexible routing (direct, topic, headers). In high-load projects we use Kafka; for low-latency scenarios, RabbitMQ. The choice also depends on whether guaranteed delivery (Kafka) or simpler setup (RabbitMQ) is needed.

How to Monitor Stuck Sagas?

Without visibility into Saga states, debugging distributed transactions is extremely difficult. We use SQL queries and alerts. We set up Prometheus + Grafana with Telegram notifications — average incident response time dropped to 5 minutes.

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 on stuck Sagas:

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

What Are the Risks When Implementing Saga?

  • Message loss — if the broker provides at-least-once guarantees, handlers must be idempotent. Otherwise duplicates can occur: in one project this led to 12% order errors.
  • Inconsistency — due to compensation failures, data may remain in an intermediate state. We ensure exactly one compensation call and its idempotency.
  • Testing complexity — requires simulating failures of each service. We deploy a staging environment with chaos engineering (e.g., Chaos Mesh).

We cover these risks with tests and monitoring. As a result, over 95% of Sagas complete successfully on the first attempt.

Deliverables and What’s Included

  1. Audit of current architecture and broker selection (Apache Kafka or RabbitMQ).
  2. Saga design: step chain, compensations, message format.
  3. Orchestrator implementation (Java/Spring) with idempotency and recovery.
  4. Broker integration.
  5. Development of compensating methods in each service.
  6. Testing failure scenarios (simulating each service crash).
  7. Monitoring setup (dashboards, alerts, manual management).
  8. Documentation and team training.
  9. One month of post-deployment support.

Our typical implementation costs $15,000–$20,000.

Estimated Timeline

Stage Duration
Design and analysis 1–2 days
Orchestrator development 2–3 days
Service integration 1–2 days
Testing and debugging 1–2 days
Monitoring and documentation 1 day
Total 7–10 days

Contact us for a consultation — we’ll analyze your architecture and suggest the best solution. Order implementation and we’ll help avoid typical pitfalls.

Get reliable distributed transactions without locks.