第 6 章 · Spring Cloud Stream 与 RabbitMQ / Kafka 概念
本章目标:理解事件驱动在微服务中的角色;掌握 Spring Cloud Stream 的 Binder、Binding、Channel 抽象;对比 RabbitMQ 与 Apache Kafka 适用场景;在 svc-spring-demo 实现「订单创建 → 异步通知库存/日志」教学链路;了解消息可靠性、幂等消费与死信队列概念。
学时建议:5~6 小时(含 2 小时 RabbitMQ 本地联调)
前置:spring-cloud-web ch01~ch05;spring-boot-web ch10 @Async;消息队列基本概念。
6.1 同步 Feign 的局限与事件驱动
ch04 中 order-svc 同步调用 product-svc 校验商品,适合强一致读场景。但以下需求更适合异步消息:
| 场景 | 同步 Feign 问题 | 消息方案 |
|---|---|---|
| 发邮件/短信通知 | 拉长订单接口 RT | 发布事件,通知服务订阅 |
| 扣减库存(最终一致) | 库存服务慢或不可用拖垮下单 | 发布 OrderCreated,inventory 异步扣 |
| 搜索索引更新 | 非核心路径 | 异步同步 ES |
| 削峰填谷 | 秒杀瞬时流量压垮 DB | 队列缓冲 |
同步(ch04):
order-svc ──Feign──► product-svc(查价)
异步(本章):
order-svc ──publish──► order-created ──► inventory-svc
└──► notify-svc
svc-spring-demo 为虚构练习;消息 Broker 地址使用localhost或mq.example.com占位,勿连接生产集群。
6.2 RabbitMQ 与 Kafka 概念对照
| 维度 | RabbitMQ(AMQP) | Kafka |
|---|---|---|
| 模型 | Exchange → Queue → Consumer | Topic → Partition → Consumer Group |
| 吞吐量 | 万级 QPS 足够多数业务 | 百万级日志流 |
| 路由 | 灵活(direct/topic/fanout) | 按 key 分区 |
| 消息回溯 | 消费即删(可 TTL/ DLX) | 按 offset 可回溯 |
| 运维 | Erlang,相对轻 | 3.x 起默认 KRaft 免 ZooKeeper(4.0 已移除 ZK),集群仍较重 |
| 教学默认 | 本章 RabbitMQ 为主 | 概念对比 + 配置差异 |
选型建议:
- 业务事件、任务队列 → RabbitMQ
- 行为日志、大数据管道 → Kafka
- Spring Cloud Stream 通过换 Binder 降低迁移成本
6.3 Spring Cloud Stream 核心抽象
┌─────────────────────────────────────────┐
│ order-svc 应用代码 │
│ @Bean Supplier / Consumer / Function │
└─────────────────┬───────────────────────┘
│ Binding: order-created-out
▼
┌─────────────────────────────────────────┐
│ Stream Bridge / Functional Model │
└─────────────────┬───────────────────────┘
│ RabbitMQ Binder
▼
┌─────────────────────────────────────────┐
│ RabbitMQ Exchange: order.events │
│ Queue: inventory.order.created │
└─────────────────────────────────────────┘
| 概念 | 说明 |
|---|---|
| Binder | 与 MQ 打交道的适配层(rabbit、kafka) |
| Binding | 连接应用与中间件目的地 |
| destination | 逻辑通道名,如 order-created |
| group | 消费组,竞争消费 vs 广播 |
Spring Cloud Stream 3.x+ 推荐 函数式编程模型(Supplier、Consumer、Function)。
6.4 本地启动 RabbitMQ
docker-compose.dev.yml 追加:
rabbitmq:
image: rabbitmq:3.13-management
ports:
- "5672:5672"
- "15672:15672"
environment:
RABBITMQ_DEFAULT_USER: guest
RABBITMQ_DEFAULT_PASS: guest
管理台:http://127.0.0.1:15672(guest/guest,仅本地)。
6.5 order-svc:发布订单创建事件
6.5.1 依赖
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-stream-rabbit</artifactId>
</dependency>
Kafka 则换:
<artifactId>spring-cloud-starter-stream-kafka</artifactId>
6.5.2 事件 DTO(svc-common)
@Data
@AllArgsConstructor
@NoArgsConstructor
public class OrderCreatedEvent {
private Long orderId;
private Long userId;
private Long productId;
private Integer quantity;
private BigDecimal amount;
private Instant createdAt;
}
6.5.3 配置
spring:
cloud:
stream:
bindings:
orderCreated-out-0:
destination: order.created
content-type: application/json
rabbit:
bindings:
orderCreated-out-0:
producer:
routing-key-expression: "'order.created'"
6.5.4 使用 StreamBridge 发送
@Service
@RequiredArgsConstructor
@Slf4j
public class OrderEventPublisher {
private final StreamBridge streamBridge;
private static final String BINDING = "orderCreated-out-0";
public void publishOrderCreated(OrderCreatedEvent event) {
boolean sent = streamBridge.send(BINDING, MessageBuilder
.withPayload(event)
.setHeader("eventType", "OrderCreated")
.setHeader("X-Trace-Id", MDC.get("traceId"))
.build());
if (!sent) {
log.error("Failed to publish OrderCreatedEvent orderId={}", event.getOrderId());
throw new BusinessException(ErrorCode.MESSAGE_PUBLISH_FAILED);
}
log.info("Published OrderCreatedEvent orderId={}", event.getOrderId());
}
}
在 OrderCreateService.create 成功落库后调用 publishOrderCreated。
6.6 inventory-svc(选修模块):消费事件
新建轻量 svc-inventory 或先在 order-svc 内写 Consumer 做教学:
6.6.1 Consumer 函数
@Configuration
public class OrderEventConsumerConfig {
@Bean
public Consumer<OrderCreatedEvent> orderCreated() {
return event -> {
log.info("Received OrderCreated orderId={} productId={} qty={}",
event.getOrderId(), event.getProductId(), event.getQuantity());
// 扣减库存(教学:打印日志;完整逻辑需 inventory 表)
};
}
}