Skip to content

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 不负责定义业务消息格式,而是把刷新事件包装后发送到广播目的地,让每个应用实例收到并处理一次。

mermaid
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 关注“控制事件如何广播”;两者可以共用代理,但消费组、重试、幂等和失败处理边界不同。

为什么需要消息事件

同步调用的问题:

mermaid
flowchart TD
    A["订单服务"] --> B["库存服务"]
    B --> C["积分服务"]
    C --> D["短信服务"]
    D --> E["任意服务慢都会拖慢主链路"]

事件驱动方式:

mermaid
flowchart TD
    A["订单服务创建订单"] --> B["发送 OrderCreatedEvent"]
    B --> C["消息中间件"]
    C --> D["多个独立业务消费组"]
    D --> E["库存、积分、短信和数仓分别消费"]

好处:

  • 主链路更短。
  • 下游服务互相解耦。
  • 峰值流量可以缓冲。
  • 新增订阅方不用改订单服务。

代价:

  • 数据通常是最终一致。
  • 消费端必须幂等。
  • 需要处理重试、死信、积压。
  • 调用链排查更复杂,需要 traceId 进入消息。

Spring Cloud Stream 是什么

Spring Cloud Stream 是消息编程模型抽象。业务代码面对的是输入/输出通道或函数,底层通过 Binder 接到具体消息中间件。

mermaid
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消费组

函数式编程模型

当前更推荐函数式模型。

生产消息:

java
@Bean
public Supplier<OrderCreatedEvent> orderSupplier() {
    return () -> new OrderCreatedEvent(1001L, 7L);
}

消费消息:

java
@Bean
public Consumer<OrderCreatedEvent> orderConsumer() {
    return event -> {
        System.out.println("收到订单事件: " + event.getOrderId());
    };
}

事件对象:

java
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;
    }
}

配置示例:

yaml
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 是消息总线。它把多个服务实例连接到同一个消息通道,用来广播系统事件。

最常见场景是配置刷新:

mermaid
flowchart TD
    A["配置中心配置变更"] --> B["触发 /actuator/busrefresh"]
    B --> C["Spring Cloud Bus 发布刷新事件"]
    C --> D["消息中间件"]
    D --> E["各在线目标实例分别接收"]
    E --> F["每个实例刷新自己的配置"]

没有 Bus 时,你可能要逐个服务调用刷新接口;有 Bus 后,可以通过消息广播刷新事件。

Bus 和 Stream 的关系

对比项Spring Cloud StreamSpring Cloud Bus
定位通用消息编程模型系统事件总线
常见用途业务事件、异步处理配置刷新、服务间广播控制事件
面向对象业务消息Spring Cloud 系统事件
依赖中间件Kafka、RabbitMQ 等通常也依赖 Kafka、RabbitMQ

Bus 底层可以基于 Stream,但二者解决的问题不同。

事件驱动业务设计原则

事件命名用过去式

推荐:

text
OrderCreatedEvent
PaymentSucceededEvent
InventoryLockedEvent

不推荐:

text
CreateOrderCommand
DoPaymentMessage

事件表示事实已经发生,命令表示要求别人做事。事件驱动系统里,发布者最好只发布事实,不直接命令下游。

消费端必须幂等

消息可能重复。重复原因包括:

  • 消费成功但 ACK 失败。
  • 消费者宕机后重平衡。
  • 生产者重试。
  • 消息中间件至少一次投递。

幂等方式:

  • 消费日志表。
  • 业务唯一索引。
  • 状态机。
  • Redis 去重。

其中最容易出事故的是“业务事务”和“Broker确认”的先后顺序。更深入的Kafka Offset、Rabbit ACK、业务本地事务、重复投递和消费日志设计见:消费确认、Offset/ACK与业务事务顺序

失败要有重试和死信

失败消息不能无限阻塞主队列,也不能直接丢。

mermaid
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,还要理解事件驱动的代价:最终一致、幂等、重试、死信、积压和链路追踪。

下一步: