Spring Cloud Stream 事件驱动架构
Spring Cloud Stream 统一消息中间件编程模型:同一套 API 对接 Kafka、RabbitMQ、RocketMQ。本文讲清 Binder 抽象、通道模型与两种编程风格。
架构总览
应用代码(业务逻辑)
│
├─ @StreamListener / Function 注解
│
▼
绑定通道(Input/Output)
│
▼
Binder(绑定器)← 与具体中间件对接
│
├── Kafka Binder
├── RabbitMQ Binder
└── RocketMQ Binder生产:业务方法 → Output 通道 → Binder → MQ
消费:MQ → Binder → Input 通道 → 业务方法Binder 抽象:核心设计
一句话理解
Binder = 通道 ↔ 消息中间件 的适配层。应用代码只面向通道编程,切换中间件只换依赖与配置,代码零改动。
┌─────────────┐
应用代码 ──通道──▶ │ Binder │ ──▶ Kafka
│(抽象) │ ──▶ RabbitMQ
└─────────────┘ ──▶ RocketMQBinder SPI
java
// org.springframework.cloud.stream.binder.Binder
public interface Binder<T, C extends ConsumerProperties, P extends ProducerProperties> {
// 绑定消费者:中间件资源 → 应用通道
Binding<T> bindConsumer(String name, String group, T inboundBindTarget,
C consumerProperties);
// 绑定生产者:应用通道 → 中间件资源
Binding<T> bindProducer(String name, T outboundBindTarget,
P producerProperties);
}name:Topic/Queue 名
group:消费者组(同组竞争消费)
inboundBindTarget:应用侧的 Input 通道
outboundBindTarget:应用侧的 Output 通道绑定过程
应用启动
│
├─ 检测到 @EnableBinding(接口) / Function 定义
│
├─ 解析出绑定目标(Input/Output 通道)
│
├─ BinderFactory 创建对应 Binder(按依赖自动选择 Kafka/Rabbit)
│
├─ Binder.bindConsumer/bindProducer
│ ├─ 连接 MQ(Topic/Queue/Exchange)
│ ├─ 建立消息流
│ └─ 绑定到应用通道
│
└─ 通道可用:消息开始流动消息通道模型
三种通道
| 通道 | 接口 | 特性 |
|---|---|---|
| MessageChannel | MessageChannel | 基本发送通道 |
| SubscribableChannel | SubscribableChannel | 可被消费者订阅 |
| PollableChannel | PollableChannel | 可轮询接收 |
java
// Spring Integration 的 MessageChannel
public interface MessageChannel {
boolean send(Message<?> message);
}
// 可订阅:多个消费者注册
public interface SubscribableChannel extends MessageChannel {
boolean subscribe(MessageHandler handler);
boolean unsubscribe(MessageHandler handler);
}Message 结构
java
public interface Message<T> {
T getPayload(); // 业务数据
MessageHeaders getHeaders(); // 消息头(类型/分区/来源...)
}Message {
payload: {"orderId": 1001, "amount": 299.0}
headers: {
contentType: application/json
deliveryAttempt: 1
partition: 0
}
}注解式编程模型(@StreamListener)
绑定接口
java
// 1. 定义绑定接口
public interface OrderChannel {
@Input("order-input") // 消费:MQ → 应用
SubscribableChannel input();
@Output("order-output") // 生产:应用 → MQ
MessageChannel output();
}java
// 2. 启用绑定
@SpringBootApplication
@EnableBinding(OrderChannel.class)
public class OrderApplication { ... }消费
java
// 3. 消费消息
@Component
public class OrderConsumer {
@StreamListener("order-input")
public void onOrder(OrderMessage message) {
// 处理订单消息
log.info("收到订单: {}", message.getOrderId());
}
}生产
java
@Component
public class OrderProducer {
@Autowired
private OrderChannel channel;
public void sendOrder(OrderMessage message) {
channel.output().send(MessageBuilder.withPayload(message).build());
}
}注解参数解析
java
@StreamListener("order-input")
public void onOrder(
OrderMessage payload, // 消息体(反序列化)
@Header("contentType") String type, // 消息头
@Headers Map<String, Object> headers) {
...
}条件消费
java
// 按条件路由
@StreamListener(value = "order-input", condition = "headers['type']=='create'")
public void onCreate(OrderMessage msg) { ... }
@StreamListener(value = "order-input", condition = "headers['type']=='cancel'")
public void onCancel(OrderMessage msg) { ... }函数式编程模型(推荐)
Spring Cloud Stream 3.x 起推荐函数式模型,用 java.util.function 定义通道:
函数类型
| 函数 | 角色 | 数量 |
|---|---|---|
Consumer<T> | 消费(输入) | 1 个输入 |
Function<T, R> | 处理(输入+输出) | 1 入 1 出 |
Supplier<T> | 产生(输出) | 1 个输出 |
定义
java
@Configuration
public class StreamFunctions {
// 消费:order-in 通道
@Bean
public Consumer<OrderMessage> orderConsumer() {
return msg -> log.info("收到订单: {}", msg.getOrderId());
}
// 处理:upper-in → upper-out
@Bean
public Function<String, String> uppercase() {
return s -> s.toUpperCase();
}
// 产生:ticker-out 定时输出
@Bean
public Supplier<String> ticker() {
return () -> "tick-" + System.currentTimeMillis();
}
}绑定配置
yaml
spring:
cloud:
stream:
bindings:
orderConsumer-in-0: # <函数名>-<in/out>-<序号>
destination: order.topic # 对应 MQ Topic
group: order-group # 消费者组
uppercase-in-0:
destination: source.topic
uppercase-out-0:
destination: target.topic命名规则:<functionName>-<in|out>-<index>,多个输入/输出用序号区分。
与注解式对比
| 对比 | @StreamListener | 函数式 |
|---|---|---|
| 版本 | 传统(3.x 前) | 推荐(3.x+) |
| 通道绑定 | @EnableBinding 接口 | 函数自动绑定 |
| 测试 | 需 mock 通道 | 直接测函数 |
| 灵活性 | 注解参数丰富 | 纯函数,复用性高 |
| 声明式 | 配置多 | 配置简单 |
消息转换与类型
默认序列化
发送:对象 → JSON(Jackson)→ 消息
接收:JSON → 对象(根据方法参数类型)yaml
spring:
cloud:
stream:
bindings:
orderConsumer-in-0:
content-type: application/json # 指定内容类型自定义转换器
java
@Bean
public MessageConverter customConverter() {
return new MessageConverter() {
@Override
public Object fromMessage(Message<?> message, Class<?> targetClass) {
// 自定义反序列化
}
@Override
public Message<?> toMessage(Object payload, MessageHeaders headers) {
// 自定义序列化
}
};
}消费者组与分区
消费者组
同组消费者竞争消费(一条消息只被组内一个消费)
不同组各自全量消费(发布-订阅)yaml
spring:
cloud:
stream:
bindings:
orderConsumer-in-0:
destination: order.topic
group: group-a # 组名分区(有序处理)
生产者按 key 分区 → 同一 key 的消息进同一分区 → 有序消费yaml
# 生产者
spring:
cloud:
stream:
bindings:
uppercase-out-0:
producer:
partition-key-expression: payload.orderId # 分区键
partition-count: 3 # 分区数
# 消费者(需与生产者分区数一致)
orderConsumer-in-0:
consumer:
partitioned: true完整生产示例
yaml
spring:
cloud:
stream:
bindings:
orderConsumer-in-0:
destination: order.topic
group: order-group
consumer:
max-attempts: 3 # 消费重试次数
retry-backoff: 1000 # 重试间隔
orderSupplier-out-0:
destination: event.topic
producer:
required-groups: order-group # 持久化
kafka:
binder:
brokers: localhost:9092
auto-create-topics: truejava
@Configuration
public class OrderStream {
@Bean
public Consumer<OrderMessage> orderConsumer() {
return msg -> {
try {
// 业务处理
processOrder(msg);
} catch (Exception e) {
// 处理失败(可投递死信)
throw new RuntimeException(e);
}
};
}
}常见问题
- 切换中间件要改代码吗? 不用。换依赖(kafka binder → rabbit binder)与配置即可,业务代码不变。
- @StreamListener 不生效? 3.x 起推荐函数式,注解式需保留 @EnableBinding。
- 消息消费重试/死信? 配置 max-attempts 与死信通道(Kafka 的 DLT、Rabbit 的 DLQ)。
- 分区顺序如何保证? 分区键表达式一致 + 分区数一致,消息按 key 路由到同一分区。