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

レッスン 23: メッセージ キュー — Kafka と RabbitMQ

メッセージ キューを備えたイベント駆動型アーキテクチャ。 Apache Kafka — プロデューサー、コンシューマー、トピック。 RabbitMQ — 交換、キュー、バインディング。分散トランザクションのサーガ パターン。

💻 プログラミング — レッスン 22 レッスン 23: メッセージ キュー — Kafka と RabbitMQ

Spring Boot 4: 基本から上級まで

パート 6: マイクロサービスとプロダクション

xdev.asia

はじめに

マイクロサービスでは、密結合を回避して復元力を高めるために、サービスは非同期で通信する必要があります。 Apache Kafka と RabbitMQ は、最も人気のある 2 つのメッセージ ブローカーです。この記事では、両方を Spring Boot 4 と統合する方法について説明します。


1. イベント駆動型アーキテクチャ

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

利点:

  • デカップリング: プロデューサーは、どのコンシューマーが聞いているかを知りません。
  • 回復力: コンシューマーがダウン → メッセージがキューで待機、データ損失なし
  • スケーラビリティ: 並列処理用のコンシューマを追加します。

2. Apache Kafka

2.1 セットアップ

// 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 プロデューサー

@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 コンシューマ

@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 トピックの構成

@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 セットアップ

// 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 交換、キュー、バインディング

@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 プロデューサー

@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 コンシューマー

@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. カフカ vs RabbitMQ

特長カフカラビットMQ
モデルログベース (追加のみ)キューベース (メッセージ消費)
スループット非常に高い (数百万/秒)高 (数万/秒)
メッセージの再生はい (消費者はもう一度読みます)いいえ (ACK 後)
注文パーティション内キュー中
使用例イベントストリーミング、分析タスクキュー、RPC
複雑さ曹操中程度

5. Saga パターン — 分散トランザクション

┌────────────┐    ┌────────────┐    ┌────────────┐
│ 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());
    }
}

概要

  • イベント駆動型アーキテクチャにより、マイクロサービスの非同期通信が可能になり、復元力とスケーラビリティが向上します。
  • Kafka: ログベース、メッセージの再生、非常に高いスループット - イベント ストリーミングに適しています
  • RabbitMQ: キューベースの柔軟なルーティング、Dead Letter Exchange — タスク処理に適しています
  • Saga パターンは、補償イベントを介して分散トランザクションを管理します

演習

  1. Docker で Kafka をセットアップし、Order → Product イベント フローを実装します。OrderCreatedEvent を発行し、在庫を削減するために消費します。
  2. RabbitMQ のセットアップ: Exchange/Queue/Binding を作成し、失敗したメッセージの Dead Letter Queue を使用した通知フローを実装します。
  3. 注文フローのサーガ パターンを実装します: 注文 → 支払い → 在庫。いずれかのステップが失敗した場合の補償アクションを処理します。