簡介
在微服務中,服務需要非同步通訊以避免緊密耦合並提高彈性。 Apache Kafka 和 RabbitMQ 是兩個最受歡迎的訊息代理程式。本文介紹如何將兩者與 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.阿帕契卡夫卡
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 卡夫卡生產者
@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 卡夫卡消費者
@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.Kafka 與 RabbitMQ
| 特色 | 卡夫卡 | 兔子MQ |
|---|---|---|
| 型號 | 基於日誌(僅附加) | 基於佇列(訊息消耗) |
| 吞吐量 | 非常高(數百萬/秒) | 高(數萬/秒) |
| 留言重播 | 是(消費者再次閱讀) | 否(確認後) |
| 訂購 | 在分區 | 排隊中 |
| 使用案例 | 事件流、分析 | 任務佇列、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:基於佇列、彈性路由、死信交換-適合任務處理
- Saga模式透過補償事件管理分散式事務
練習
- 使用Docker搭建Kafka,實現訂單→產品事件流程:發布OrderCreatedEvent,消費減少庫存
- 設定RabbitMQ:建立Exchange/Queue/Binding,使用死信佇列實作失敗訊息的通知流 3.訂單流程實現Saga模式:訂單→付款→庫存,任一步驟失敗時處理補償動作