Spring Cloud Stream 消息驱动的微服务
导学
飞翔科技的订单系统最近遇到了一个问题。孔蓝提了需求:"用户下单后,要同时做三件事——扣库存、发短信通知、写审计日志。现在这三个操作串行执行,用户等 3 秒才能看到'下单成功'。"
白歌画了一张对比图:
"用消息队列解耦——这就是 Spring Cloud Stream 的价值。"白歌说。
定位与问题场景
Spring Cloud Stream 是一个消息中间件抽象层,让开发者面向统一的消息编程模型编写代码,无需关心底层是 Kafka、RabbitMQ 还是 RocketMQ。
核心概念
| 概念 | 说明 |
|---|---|
| Binder | 与消息中间件的适配层,屏蔽 RabbitMQ/Kafka 差异 |
| Destination | 消息目标(Kafka 中的 Topic,RabbitMQ 中的 Exchange) |
| Supplier | 消息生产者(函数式),输出消息 |
| Function | 消息处理器,接收消息 + 处理后输出 |
| Consumer | 消息消费者(函数式),接收消息 |
| StreamBridge | 程序化发送消息的工具类 |
Stream 3.x 函数式编程模型
Spring Cloud Stream 3.x 废弃了 @EnableBinding / @Input / @Output 注解式编程,全面转向函数式编程:
完整示例:飞翔科技订单异步处理
场景描述
用户下单后:
- 订单服务发送
OrderCreatedEvent到消息队列 - 库存服务消费事件扣减库存
- 通知服务消费事件发送短信
- 审计服务消费事件记录日志
操作前后对比:
| 维度 | 引入 Stream 前 | 引入 Stream 后 |
|---|---|---|
| 下单响应时间 | 3 秒(串行) | 0.5 秒(异步) |
| 耦合度 | 订单服务直接调用库存/通知/审计 | 订单服务只发消息,完全解耦 |
| 削峰填谷 | 无 | 消息队列缓冲 |
示例一:订单服务——发送消息(Supplier + StreamBridge)
依赖:
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-stream-binder-rabbit</artifactId>
</dependency>
配置:
spring:
cloud:
stream:
bindings:
orderCreated-out-0: # 绑定名:<functionName>-<in|out>-<index>
destination: order.created # 对应 RabbitMQ Exchange / Kafka Topic
content-type: application/json
orderRefunded-out-0:
destination: order.refunded
content-type: application/json
rabbitmq:
host: localhost
port: 5672
username: admin
password: admin
订单服务代码:
@RestController
public class OrderController {
@Autowired
private StreamBridge streamBridge;
@PostMapping("/orders")
public Order createOrder(@RequestBody OrderRequest request) {
// 1. 保存订单到数据库(核心业务,同步)
Order order = orderService.save(request);
// 2. 发送事件(通知+审计,异步)
OrderCreatedEvent event = new OrderCreatedEvent(
order.getId(), order.getUserId(), order.getTotalAmount()
);
streamBridge.send("orderCreated-out-0", event);
// 3. 立即返回(不等待库存扣减、短信、日志)
return order;
}
}
示例二:库存服务——消费消息(Consumer)
@Configuration
public class InventoryConsumer {
@Bean
public Consumer<OrderCreatedEvent> orderCreated() {
return event -> {
log.info("收到订单创建事件,扣减库存:orderId={}", event.getOrderId());
inventoryService.deduct(event);
};
}
}
spring:
cloud:
stream:
bindings:
orderCreated-in-0:
destination: order.created # 必须与生产者 destination 一致
group: inventory-consumer-group # 消费者组(持久化 + 负载均衡)
示例三:通知服务——消费消息(Consumer)
@Configuration
public class NotificationConsumer {
@Bean
public Consumer<OrderCreatedEvent> orderCreated() {
return event -> {
log.info("收到订单创建事件,发送短信:orderId={}", event.getOrderId());
smsService.sendOrderConfirmation(event.getUserId(), event.getOrderId());
};
}
}
示例四:消息处理——Function(接收 + 转换 + 输出)
@Configuration
public class OrderProcessor {
@Bean
public Function<OrderCreatedEvent, AuditLog> orderAudit() {
return event -> {
log.info("审计订单:orderId={}", event.getOrderId());
return new AuditLog("ORDER_CREATED", event.getOrderId(), Instant.now());
};
}
}
spring:
cloud:
stream:
bindings:
orderAudit-in-0:
destination: order.created
orderAudit-out-0:
destination: audit.log
消费组与分区
消费组(Consumer Group):
# 库存服务部署 3 个实例,同一 group 内只有一个实例消费
spring.cloud.stream.bindings.orderCreated-in-0.group: inventory-consumer-group
# 通知服务部署 2 个实例
spring.cloud.stream.bindings.orderCreated-in-0.group: notification-consumer-group
效果:消息只会被每个 group 中的一个实例消费,不同 group 各自独立消费(类似 Kafka 的消费者组)。
分区(Partitioning):
spring:
cloud:
stream:
bindings:
orderCreated-out-0:
producer:
partition-key-expression: payload.orderId
partition-count: 3
效果:相同 orderId 的消息总是路由到同一分区,保证有序处理。
易错场景
1. 消费者组缺失导致消息丢失
# ❌ 错误:未配置 group,默认匿名消费者,重启后会丢失所有未消费消息
spring.cloud.stream.bindings.orderCreated-in-0:
destination: order.created
# ✅ 正确:配置 group,消息持久化
spring.cloud.stream.bindings.orderCreated-in-0:
destination: order.created
group: my-group
2. 函数命名与 binding 命名不匹配
@Bean
public Consumer<OrderCreatedEvent> processOrder() { ... }
// 期望的 binding 名:processOrder-in-0
// 如果配置文件写成了 orderCreated-in-0 → 绑定失败
3. 忘记配置 content-type
# 默认 content-type 为 application/json
# 但发送字节数组或纯文本时需要显式指定
spring.cloud.stream.bindings.myEvent-out-0.content-type: text/plain
4. DLQ(死信队列)未配置导致异常消息丢弃
# RabbitMQ 死信队列
spring:
cloud:
stream:
rabbit:
bindings:
orderCreated-in-0:
consumer:
auto-bind-dlq: true
dlq-ttl: 5000
面试考点
Spring Cloud Stream 的 Binder 是什么?
Binder 是 Spring Cloud Stream 与底层消息中间件的适配层,负责将消息通道(Destination)映射为具体的中间件概念(Kafka 的 Topic、RabbitMQ 的 Exchange/Queue)。开发者面向统一的
Supplier/Function/Consumer编程,切换中间件只需更换 Binder 依赖,无需修改业务代码。
Consumer Group 的原理和作用?
相同的 group 内的多个实例共享消息(竞争消费,每条消息只被一个实例处理),实现负载均衡。不同的 group 各自独立消费全量消息(广播模式)。group 还提供持久化功能:消息中间件会记录每个 group 的消费偏移量,服务重启后从上次位置继续消费。不配置 group 的消费者是匿名消费者,重启后消息会丢失。
Stream 3.x 为什么废弃 @EnableBinding?
Spring Cloud Stream 3.x 全面转向函数式编程模型,原因:
- 注解式编程(
@Input/@Output/@StreamListener)需要维护繁琐的接口声明- 函数式 Bean(
Supplier/Function/Consumer)更符合 Spring Boot 的 Auto-Configuration 理念- 框架可以根据函数的输入输出类型自动推断消息转换逻辑
- 更易于单元测试
StreamBridge 与 Supplier 的区别?
Supplier:声明式,框架按配置的轮询间隔自动调用,适合轮询数据库、定时任务等场景StreamBridge:程序化,业务代码主动调用.send()发送消息,适合用户操作触发的场景(如 HTTP 请求 → 发消息)- 大多数业务场景使用
StreamBridge,定时或流式数据源使用Supplier
小结
Spring Cloud Stream 通过 Binder 抽象层屏蔽了消息中间件的差异,函数式编程模型让消息处理代码更简洁。消费组实现了消费端的负载均衡和持久化,StreamBridge 让业务代码可以灵活地发送消息。至此,微服务的同步调用链路和异步消息通道都已打通。下一章进入可观测性领域——分布式追踪。