Introduction
In microservices, services need to communicate asynchronously to avoid tight coupling and increase resilience. Apache Kafka and RabbitMQ are the two most popular message brokers. This article explains how to integrate both with Spring Boot 4.
1. Event-Driven Architecture
Synchronous (REST):
Order Service ──HTTP──► Product Service (tight coupling)
Asynchronous (Message Queue):
Order Service ──publish──► Message Broker ──consume──► Product Service
──consume──► Notification Service
──consume──► Analytics Service
Benefits:
- Decoupling: Producer does not know which consumer is listening
- Resilience: Consumer down → messages wait in queue, no data loss
- Scalability: Add consumers for parallel processing
2. Apache Kafka
2.1 Setup
// build.gradle.kts
implementation("org.springframework.kafka:spring-kafka")
# application.yml
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
consumer:
group-id: ${spring.application.name}
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "com.example.*"
2.2 Kafka Producer
@Service
public class OrderEventProducer {
private final KafkaTemplate<String, Object> kafkaTemplate;
public OrderEventProducer(KafkaTemplate<String, Object> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void publishOrderCreated(OrderCreatedEvent event) {
kafkaTemplate.send("order-events", event.orderId().toString(), event)
.whenComplete((result, ex) -> {
if (ex != null) {
log.error("Failed to publish event: {}", ex.getMessage());
} else {
log.info("Published to partition {} offset {}",
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
}
});
}
}
2.3 Kafka Consumer
@Component
public class OrderEventConsumer {
@KafkaListener(topics = "order-events", groupId = "product-service")
public void handleOrderCreated(OrderCreatedEvent event) {
log.info("Received order event: {}", event.orderId());
inventoryService.decreaseStock(event.items());
}
// Consumer với manual acknowledgment
@KafkaListener(topics = "payment-events", groupId = "order-service",
containerFactory = "kafkaManualAckListenerContainerFactory")
public void handlePaymentResult(PaymentResultEvent event,
Acknowledgment ack) {
try {
orderService.updatePaymentStatus(event);
ack.acknowledge(); // Commit offset sau khi xử lý thành công
} catch (Exception e) {
log.error("Failed to process payment event", e);
// Không ack → Kafka sẽ retry
}
}
}
2.4 Kafka Topics Configuration
@Configuration
public class KafkaTopicConfig {
@Bean
public NewTopic orderEventsTopic() {
return TopicBuilder.name("order-events")
.partitions(3)
.replicas(1)
.config(TopicConfig.RETENTION_MS_CONFIG, "604800000") // 7 days
.build();
}
@Bean
public NewTopic paymentEventsTopic() {
return TopicBuilder.name("payment-events")
.partitions(3)
.replicas(1)
.build();
}
}
3. RabbitMQ
3.1 Setup
// build.gradle.kts
implementation("org.springframework.boot:spring-boot-starter-amqp")
# application.yml
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
3.2 Exchange, Queue, Binding
@Configuration
public class RabbitMQConfig {
public static final String ORDER_EXCHANGE = "order.exchange";
public static final String ORDER_CREATED_QUEUE = "order.created.queue";
public static final String ORDER_CREATED_KEY = "order.created";
@Bean
public TopicExchange orderExchange() {
return new TopicExchange(ORDER_EXCHANGE);
}
@Bean
public Queue orderCreatedQueue() {
return QueueBuilder.durable(ORDER_CREATED_QUEUE)
.withArgument("x-dead-letter-exchange", "dlx.exchange")
.withArgument("x-dead-letter-routing-key", "dlx.order.created")
.build();
}
@Bean
public Binding orderCreatedBinding() {
return BindingBuilder
.bind(orderCreatedQueue())
.to(orderExchange())
.with(ORDER_CREATED_KEY);
}
@Bean
public Jackson2JsonMessageConverter messageConverter() {
return new Jackson2JsonMessageConverter();
}
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory factory,
Jackson2JsonMessageConverter converter) {
RabbitTemplate template = new RabbitTemplate(factory);
template.setMessageConverter(converter);
return template;
}
}
3.3 RabbitMQ Producer
@Service
public class OrderEventPublisher {
private final RabbitTemplate rabbitTemplate;
public OrderEventPublisher(RabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;
}
public void publishOrderCreated(OrderCreatedEvent event) {
rabbitTemplate.convertAndSend(
RabbitMQConfig.ORDER_EXCHANGE,
RabbitMQConfig.ORDER_CREATED_KEY,
event
);
log.info("Published order event: {}", event.orderId());
}
}
3.4 RabbitMQ Consumer
@Component
public class OrderEventHandler {
@RabbitListener(queues = RabbitMQConfig.ORDER_CREATED_QUEUE)
public void handleOrderCreated(OrderCreatedEvent event) {
log.info("Processing order: {}", event.orderId());
inventoryService.decreaseStock(event.items());
}
}
4. Kafka vs RabbitMQ
| Features | Kafka | RabbitMQ |
|---|---|---|
| Model | Log-based (append-only) | Queue-based (message consumed) |
| Throughput | Very high (millions/sec) | High (tens of thousands/sec) |
| Message Replay | Yes (consumer reads again) | No (after ack) |
| Ordering | In partition | In queue |
| Use cases | Event streaming, analytics | Task queue, RPC |
| Complexity | Cao | Moderate |
5. Saga Pattern — Distributed Transactions
┌────────────┐ ┌────────────┐ ┌────────────┐
│ Order │ │ Payment │ │ Inventory │
│ Service │ │ Service │ │ Service │
│ │ │ │ │ │
│ 1.Create │───►│ 2.Charge │───►│ 3.Reserve │
│ Order │ │ Payment │ │ Stock │
│ │ │ │ │ │
│ 6.Confirm │◄───│ 5.Confirm │◄───│ 4.Confirm │
│ /Cancel │ │ /Refund │ │ /Release │
└────────────┘ └────────────┘ └────────────┘
// Choreography-based Saga
@Component
public class OrderSagaHandler {
@KafkaListener(topics = "payment-completed")
public void onPaymentCompleted(PaymentCompletedEvent event) {
// Step: reserve inventory
inventoryService.reserve(event.orderId(), event.items());
kafkaTemplate.send("inventory-reserved",
new InventoryReservedEvent(event.orderId()));
}
@KafkaListener(topics = "payment-failed")
public void onPaymentFailed(PaymentFailedEvent event) {
// Compensating action: cancel order
orderService.cancelOrder(event.orderId());
}
@KafkaListener(topics = "inventory-failed")
public void onInventoryFailed(InventoryFailedEvent event) {
// Compensating actions: refund + cancel
paymentService.refund(event.orderId());
orderService.cancelOrder(event.orderId());
}
}
Summary
- Event-Driven Architecture helps microservices communicate asynchronously, increasing resilience and scalability
- Kafka: log-based, replay messages, extremely high throughput — suitable for event streaming
- RabbitMQ: queue-based, flexible routing, Dead Letter Exchange — suitable for task processing
- Saga pattern manages distributed transactions via compensating events
Exercises
- Setup Kafka with Docker, implement Order → Product event flow: publish OrderCreatedEvent, consume to reduce stock
- Setup RabbitMQ: create Exchange/Queue/Binding, implement notification flow with Dead Letter Queue for failed messages
- Implement Saga pattern for Order flow: Order → Payment → Inventory, handle compensating actions when any step fails