Spring Cloud Stream:消息驱动架构与 Binder
前面章节解决的是同步调用问题:请求-响应、你在等我。很多业务其实不需要同步等结果:下单后发通知、订单状态变更后触发后续处理,同步等待只会拉长响应时间、放大故障耦合。消息驱动把"动作"变成"事件"异步流转,Spring Cloud Stream 提供了与具体 MQ 解耦的统一编程模型。本章讲清概念,第 20 章用 RabbitMQ 落地。
直接使用 MQ 客户端的问题
RabbitTemplate 或 KafkaProducer 本身并不难用,难的是换厂商:代码里到处是 rabbit 的 API,某天要换 Kafka 就得大改。Spring Cloud Stream 的思路是再加一层抽象:业务代码只面向"发消息/收消息"的通用模型,具体对接哪个 MQ 由 Binder 实现,切换 MQ 时业务代码基本不动。
三个核心概念
目的地(destination):消息的逻辑目标名,业务眼里是"order-events 这个事件",Binder 负责把它映射成 RabbitMQ 的交换机/队列或 Kafka 的 topic。绑定(binding):把业务函数与某个 destination 连接起来的配置单元,形如"某个函数进/出某个 destination"。Binder:连接抽象层与具体 MQ 的适配器,RabbitMQ 有 Rabbit Binder,Kafka 有 Kafka Binder,Stream 通过引入哪个 binder 依赖决定用哪个 MQ。
函数式编程模型
早期 Stream 用 @EnableBinding、@Input、@Output 注解声明通道,这套注解在后续版本逐步废弃;新版推荐纯函数式模型——用 Spring 标准的函数 Bean 表达消息收发:
- Supplier:无入参、有出参,表示消息源(生产消息),绑定形如 函数名-out-0;
- Consumer:有入参、无出参,表示消息去向(消费消息),绑定形如 函数名-in-0;
- Function:有入参有出参,表示"接收→处理→产出"的转换。
框架根据函数 Bean 自动推断输入输出,开发者只需要在配置里用 spring.cloud.function.definition 显式声明启用哪些函数。
消费组与广播
多个消费者实例订阅同一 destination 时,消息该给谁?答案是消费组(group):同一组内的实例分摊消息(每条消息只被组内一个实例消费,实现负载均衡);不同组的消费者各自都收到完整消息(广播语义)。组内某实例挂掉,消息由组内其他实例继续消费,这是消息可靠性的基础。RabbitMQ 场景下,每个"destination + group"对应一个持久队列。
配置长什么样
一个最小配置只包含函数定义与绑定声明,例如声明"orderSupplier 函数发送到 order-events 目的地":
spring:
cloud:
stream:
bindings:
orderSupplier-out-0:
destination: order-events
function:
definition: orderSupplier
绑定名(orderSupplier-out-0)的规律是 函数名 + in/out + 序号,序号 0 表示第一个输入/输出,多输入输出时递增。
何时引入 Stream
消息场景都值得用:发件人只需要"把事件发出去",不关心下游是谁、有几个;Binder 抽象让测试与切换更从容(Stream 官方提供测试用 Binder 以支撑不依赖真实 MQ 的集成测试,具体坐标以官方文档为准)。若团队明确"只用 RabbitMQ 且永不更换",直接用原生客户端也可以——Stream 的价值在抽象与规范,不在性能。
消息模型设计建议
在写第一个事件之前,先立几条约定,收益会随着服务数量增长越来越明显:
- 事件命名用领域动词的过去式(如 OrderCreated),表达"已经发生的事实"而非"去做什么的命令";
- 事件内容自包含:带上订单号、金额等关键字段,消费者尽量不回头查发送方接口,降低耦合;
- 演进只增不改:字段只增不删,消费者容忍未知字段;破坏性变更用新事件代替修改旧事件;
- 消费幂等:消息可能重复投递,消费侧要按业务键去重,不能假设"每条只来一次"。
小结
Stream 用"destination/binding/Binder"三层把消息编程从具体 MQ 中解放出来,函数式模型(Supplier/Consumer/Function)+ definition 声明是当前主流写法,消费组解决负载均衡与广播语义。概念就绪,下一章把 RabbitMQ 真实跑起来收发消息。