Event-Driven Microservices with Apache Kafka, Redis Caching and Transactional Outbox Pattern
When you scale microservices under high traffic, traditional HTTP REST calls between services break down quickly. Cascading timeouts, network latency spikes and database deadlocks will crash your order processing pipeline during traffic bursts.
I experienced this firsthand while building high-concurrency backend pipelines and optimizing infrastructure costs to under ₹100 per month (using Koyeb, Vercel and optimized DB pooling). That same low-cost, high-performance architecture enabled our e-commerce platform to generate ₹1 Lakh+ in revenue in just 2 months without single-point-of-failure downtimes.
Here is the exact production-grade blueprint for building event-driven Spring Boot microservices with Apache Kafka, Redis caching and the Transactional Outbox Pattern.
1. The Dual-Write Problem in Microservices
Suppose your order-service needs to save a new order to PostgreSQL and notify payment-service and inventory-service via Kafka.
If you write code like this:
@Transactional
public OrderResponse createOrder(CreateOrderRequest request) {
Order order = orderRepository.save(new Order(request));
kafkaTemplate.send("order-created-topic", new OrderCreatedEvent(order.getId()));
return orderMapper.toResponse(order);
}
This code has a catastrophic flaw known as the Dual-Write Problem:
- If the database transaction commits but Kafka crashes right after, your order is saved but event consumers never receive it (data drift).
- If Kafka succeeds but the database transaction rolls back due to a constraint violation, an event is published for an order that never existed.
2. Solving Dual-Write with the Transactional Outbox Pattern
Instead of publishing to Kafka inside the database transaction, write an outbox record inside the same ACID transaction as your entity.
Outbox Entity
@Entity
@Table(name = "outbox_events")
public class OutboxEvent {
@Id
@GeneratedValue(strategy = GenerationType.UUID)
private UUID id;
@Column(nullable = false)
private String aggregateType;
@Column(nullable = false)
private String aggregateId;
@Column(nullable = false)
private String eventType;
@Column(columnDefinition = "TEXT", nullable = false)
private String payload;
@Column(nullable = false)
private Instant createdAt;
@Enumerated(EnumType.STRING)
private OutboxStatus status = OutboxStatus.PENDING;
// Constructors, getters and setters
}
Transactional Order Service Implementation
@Service
public class OrderService {
private final OrderRepository orderRepository;
private final OutboxRepository outboxRepository;
private final ObjectMapper objectMapper;
public OrderService(OrderRepository orderRepository,
OutboxRepository outboxRepository,
ObjectMapper objectMapper) {
this.orderRepository = orderRepository;
this.outboxRepository = outboxRepository;
this.objectMapper = objectMapper;
}
@Transactional
public OrderResponse createOrder(CreateOrderRequest request) {
Order order = orderRepository.save(new Order(request.customerCode(), request.totalAmount()));
OutboxEvent outbox = new OutboxEvent();
outbox.setAggregateType("Order");
outbox.setAggregateId(order.getId().toString());
outbox.setEventType("ORDER_CREATED");
outbox.setPayload(objectMapper.writeValueAsString(new OrderCreatedEvent(order.getId(), order.getTotalAmount())));
outbox.setCreatedAt(Instant.now());
outboxRepository.save(outbox);
return new OrderResponse(order.getId(), "PENDING");
}
}
3. High-Throughput Outbox Publisher with Spring Scheduler & Kafka
Now create a dedicated background publisher using Spring Virtual Threads (spring.threads.virtual.enabled=true) to process outbox records and publish to Apache Kafka:
@Component
public class OutboxPublisher {
private static final Logger log = LoggerFactory.getLogger(OutboxPublisher.class);
private final OutboxRepository outboxRepository;
private final KafkaTemplate<String, String> kafkaTemplate;
public OutboxPublisher(OutboxRepository outboxRepository, KafkaTemplate<String, String> kafkaTemplate) {
this.outboxRepository = outboxRepository;
this.kafkaTemplate = kafkaTemplate;
}
@Scheduled(fixedDelay = 500)
@Transactional
public void publishPendingEvents() {
List<OutboxEvent> pending = outboxRepository.findTop50ByStatusOrderByCreatedAtAsc(OutboxStatus.PENDING);
for (OutboxEvent event : pending) {
kafkaTemplate.send("order-events", event.getAggregateId(), event.getPayload())
.whenComplete((result, ex) -> {
if (ex == null) {
event.setStatus(OutboxStatus.PROCESSED);
outboxRepository.save(event);
} else {
log.error("Failed to publish outbox event id: {}", event.getId(), ex);
}
});
}
}
}
4. Idempotent Consumer & Redis Distributed Locking
Network retries in Kafka mean your consumers WILL receive duplicate events. You must enforce idempotency at the consumer layer using Redis:
@Component
public class PaymentOrderConsumer {
private static final Logger log = LoggerFactory.getLogger(PaymentOrderConsumer.class);
private final StringRedisTemplate redisTemplate;
public PaymentOrderConsumer(StringRedisTemplate redisTemplate) {
this.redisTemplate = redisTemplate;
}
@KafkaListener(topics = "order-events", groupId = "payment-service-group")
public void consumeOrderEvent(ConsumerRecord<String, String> record) {
String eventId = record.key();
String lockKey = "idempotency:order:" + eventId;
Boolean acquired = redisTemplate.opsForValue().setIfAbsent(lockKey, "PROCESSED", Duration.ofHours(24));
if (Boolean.FALSE.equals(acquired)) {
log.info("Duplicate event detected for order id: {}. Skipping execution.", eventId);
return;
}
// Process payment processing logic safely
log.info("Successfully processed payment for order id: {}", eventId);
}
}
Key Performance Results
By combining Virtual Threads, Redis Idempotency Keys and the Transactional Outbox Pattern:
- Zero Data Loss: Guaranteed atomic persistence of business state and event logs.
- Sub-10ms P99 Latency: REST responses return immediately after database write without waiting for external Kafka roundtrips.
- ₹100/month Infra Cost: Lightweight memory footprint allowing full stack deployment on free/low-cost cloud tiers.
Top comments (0)