Spring Cloud Stream与Bus
微服务之间不一定都要同步 HTTP 调用。很多业务更适合通过消息事件解耦,例如订单创建后通知库存、积分、风控、搜索、数据仓库。
本页负责从零建立概念和基本用法。启动绑定源码、Binder选择、消费组、分区、重试、Kafka/Rabbit差异、积压治理、Bus广播失败窗口和双版本可运行Demo,请继续学习Stream与Bus内部原理与生产治理;面试复习单独进入Stream与Bus独立面试题。
学习版本说明
| 学习线 | 推荐依赖组合 | 本专栏定位 |
|---|---|---|
| 存量系统 | JDK 8 + Spring Boot 2.7.x + Spring Cloud 2021.0.x + Stream 3.2.x | 维护旧系统、迁移旧注解模型 |
| 现代系统 | Java 17+ + Spring Boot 3.x + 与之匹配的Spring Cloud发布列车 + Stream 4.x | 新项目主线、函数式模型与现代可观测性 |
Boot、Cloud和Stream必须按官方兼容矩阵配套。不能只升级Boot,也不能把Record等Java 17语法复制到JDK 8项目。专栏在相同原理下分别标明新旧版本差异。
Spring Cloud 里和消息相关的常见技术有两个:
- Spring Cloud Stream:抽象消息中间件,让业务代码不用直接绑定 Kafka 或 RabbitMQ API。
- Spring Cloud Bus:基于消息总线广播事件,常见于配置刷新。
阅读入口:先区分业务事件和控制事件
按“业务函数 → Binder → 消息代理 → 消费组 → ACK/重试”阅读 Stream;再看“配置中心 → Bus 事件 → 多实例刷新”阅读 Bus。两条链共用代理,但消息语义、幂等和故障处理不同。
核心原理:Stream 抽象数据面,Bus 复用消息做控制面广播
Stream 的函数式模型把 Supplier/Function/Consumer 绑定到 Binder,Binder 再把目的地映射为 Kafka Topic 或 RabbitMQ Exchange。Bus 不负责定义业务消息格式,而是把刷新事件包装后发送到广播目的地,让每个应用实例收到并处理一次。
flowchart LR
A["业务函数 Consumer"] --> B["Stream Binder"]
B --> C["Kafka Topic/Rabbit Exchange"]
C --> D["消费组与分区分配"]
D --> E["业务处理、幂等与ACK"]
F["配置中心发布变更"] --> G["Bus事件"]
G --> C
C --> H["各实例Refresh监听器"]Stream 关注“业务事件如何生产和消费”,Bus 关注“控制事件如何广播”;两者可以共用代理,但消费组、重试、幂等和失败处理边界不同。
为什么需要消息事件
同步调用的问题:
flowchart TD
A["订单服务"] --> B["库存服务"]
B --> C["积分服务"]
C --> D["短信服务"]
D --> E["任意服务慢都会拖慢主链路"]事件驱动方式:
flowchart TD
A["订单服务创建订单"] --> B["发送 OrderCreatedEvent"]
B --> C["消息中间件"]
C --> D["多个独立业务消费组"]
D --> E["库存、积分、短信和数仓分别消费"]好处:
- 主链路更短。
- 下游服务互相解耦。
- 峰值流量可以缓冲。
- 新增订阅方不用改订单服务。
代价:
- 数据通常是最终一致。
- 消费端必须幂等。
- 需要处理重试、死信、积压。
- 调用链排查更复杂,需要 traceId 进入消息。
Spring Cloud Stream 是什么
Spring Cloud Stream 是消息编程模型抽象。业务代码面对的是输入/输出通道或函数,底层通过 Binder 接到具体消息中间件。
flowchart TD
A["业务代码"] --> B["Spring Cloud Stream"]
B --> C["具体Binder适配层"]
C --> D["Kafka、RabbitMQ等Broker"]核心概念:
| 概念 | 说明 |
|---|---|
| Binder | 连接具体消息中间件的适配层 |
| Binding | 业务函数和 Topic/Queue 的绑定关系 |
| Supplier | 生产消息 |
| Function | 处理输入并产生输出 |
| Consumer | 消费消息 |
| Destination | 目标 Topic 或 Queue |
| Group | 消费组 |
函数式编程模型
当前更推荐函数式模型。
生产消息:
@Bean
public Supplier<OrderCreatedEvent> orderSupplier() {
return () -> new OrderCreatedEvent(1001L, 7L);
}消费消息:
@Bean
public Consumer<OrderCreatedEvent> orderConsumer() {
return event -> {
System.out.println("收到订单事件: " + event.getOrderId());
};
}事件对象:
public class OrderCreatedEvent {
private Long orderId;
private Long userId;
public OrderCreatedEvent() {
}
public OrderCreatedEvent(Long orderId, Long userId) {
this.orderId = orderId;
this.userId = userId;
}
public Long getOrderId() {
return orderId;
}
public void setOrderId(Long orderId) {
this.orderId = orderId;
}
public Long getUserId() {
return userId;
}
public void setUserId(Long userId) {
this.userId = userId;
}
}配置示例:
spring:
cloud:
function:
definition: orderConsumer
stream:
bindings:
orderConsumer-in-0:
destination: order.created
group: search-index-service
kafka:
binder:
brokers: localhost:9092这里的重点:
destination对应 Kafka Topic 或 RabbitMQ Exchange/Queue。group对应消费组。- 同一个 group 内竞争消费,不同 group 广播消费。
Stream 和直接写 Kafka/RabbitMQ API 有什么区别
| 对比项 | Spring Cloud Stream | 直接使用 Kafka/RabbitMQ API |
|---|---|---|
| 抽象层级 | 高,业务面对统一模型 | 低,直接面对中间件细节 |
| 切换中间件 | 相对容易 | 改动较大 |
| 细节控制 | 部分高级特性需要透传配置 | 控制最完整 |
| 学习成本 | 先学 Stream 模型 | 先学具体 MQ |
| 适合场景 | 普通事件驱动、业务解耦 | 深度使用 Kafka/RabbitMQ 特性 |
如果业务大量依赖 Kafka 分区、事务、Exactly Once、复杂消费者配置,直接使用 Kafka 客户端或 Spring Kafka 可能更清晰。如果只是标准事件收发,Stream 能减少样板代码。
Spring Cloud Bus 是什么
Spring Cloud Bus 是消息总线。它把多个服务实例连接到同一个消息通道,用来广播系统事件。
最常见场景是配置刷新:
flowchart TD
A["配置中心配置变更"] --> B["触发 /actuator/busrefresh"]
B --> C["Spring Cloud Bus 发布刷新事件"]
C --> D["消息中间件"]
D --> E["各在线目标实例分别接收"]
E --> F["每个实例刷新自己的配置"]没有 Bus 时,你可能要逐个服务调用刷新接口;有 Bus 后,可以通过消息广播刷新事件。
Bus 和 Stream 的关系
| 对比项 | Spring Cloud Stream | Spring Cloud Bus |
|---|---|---|
| 定位 | 通用消息编程模型 | 系统事件总线 |
| 常见用途 | 业务事件、异步处理 | 配置刷新、服务间广播控制事件 |
| 面向对象 | 业务消息 | Spring Cloud 系统事件 |
| 依赖中间件 | Kafka、RabbitMQ 等 | 通常也依赖 Kafka、RabbitMQ |
Bus 底层可以基于 Stream,但二者解决的问题不同。
事件驱动业务设计原则
事件命名用过去式
推荐:
OrderCreatedEvent
PaymentSucceededEvent
InventoryLockedEvent不推荐:
CreateOrderCommand
DoPaymentMessage事件表示事实已经发生,命令表示要求别人做事。事件驱动系统里,发布者最好只发布事实,不直接命令下游。
消费端必须幂等
消息可能重复。重复原因包括:
- 消费成功但 ACK 失败。
- 消费者宕机后重平衡。
- 生产者重试。
- 消息中间件至少一次投递。
幂等方式:
- 消费日志表。
- 业务唯一索引。
- 状态机。
- Redis 去重。
其中最容易出事故的是“业务事务”和“Broker确认”的先后顺序。更深入的Kafka Offset、Rabbit ACK、业务本地事务、重复投递和消费日志设计见:消费确认、Offset/ACK与业务事务顺序。
失败要有重试和死信
失败消息不能无限阻塞主队列,也不能直接丢。
flowchart TD
A["消费者处理消息"] --> B["成功则提交确认"]
B --> C["失败则有限重试"]
C --> D["仍失败进入死信队列"]
D --> E["告警、修复和人工补偿"]商业常用场景
| 场景 | 用法 |
|---|---|
| 订单事件广播 | 订单服务发布事件,库存、积分、风控分别消费 |
| 配置刷新 | Bus 广播配置刷新事件 |
| 数据同步 | 业务服务发布变更事件,下游同步搜索索引 |
| 审计日志 | 业务操作发送审计事件 |
| 削峰填谷 | 高峰期消息先堆积,下游慢慢消费 |
常见问题
用了消息是不是就一定可靠
不是。可靠性取决于:
- 生产者是否确认发送成功。
- Broker 是否持久化。
- 消费者是否业务成功后确认。
- 消费端是否幂等。
- 是否有重试和死信。
Stream 能不能完全屏蔽 Kafka 和 RabbitMQ 差异
不能。Stream 能屏蔽常规收发模型,但 Kafka 的分区、offset、消费组,RabbitMQ 的 Exchange、Routing Key、ACK、DLX 等底层差异仍然会影响生产设计。
配置刷新一定要用 Bus 吗
不一定。小系统可以手动刷新或重启;大系统实例多、配置变更频繁时,Bus 的广播能力更有价值。
小结
Spring Cloud Stream 和 Bus 解决的是微服务里的异步事件问题。Stream 面向业务消息,Bus 面向系统事件广播。学习时不能只会配置 Binder,还要理解事件驱动的代价:最终一致、幂等、重试、死信、积压和链路追踪。
下一步:
