Chuyển đến nội dung chính

Lesson 23: Message Queue — Kafka & RabbitMQ

Event-Driven Architecture with Message Queues. Apache Kafka — producer, consumer, topics. RabbitMQ — exchanges, queues, bindings. Saga pattern for distributed transactions.

💻 Programming — Lesson 22 Lesson 23: Message Queue — Kafka & RabbitMQ

Spring Boot 4: From Basics to Advanced

Part 6: Microservices & Production

xdev.asia

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

FeaturesKafkaRabbitMQ
ModelLog-based (append-only)Queue-based (message consumed)
ThroughputVery high (millions/sec)High (tens of thousands/sec)
Message ReplayYes (consumer reads again)No (after ack)
OrderingIn partitionIn queue
Use casesEvent streaming, analyticsTask queue, RPC
ComplexityCaoModerate

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

  1. Setup Kafka with Docker, implement Order → Product event flow: publish OrderCreatedEvent, consume to reduce stock
  2. Setup RabbitMQ: create Exchange/Queue/Binding, implement notification flow with Dead Letter Queue for failed messages
  3. Implement Saga pattern for Order flow: Order → Payment → Inventory, handle compensating actions when any step fails