Spring Cloud Stream 源码阅读
从 @EnableBinding / 函数定义到消息真正流入流出,中间发生了什么?本文从源码拆解 Binder SPI、绑定流程、消息通道的创建与消息处理链。
核心模块与启动链路
SpringApplication 启动
│
├─ 1. 自动配置:StreamAutoConfiguration
│ ├─ BinderFactory(创建 Binder)
│ ├─ BindingService(管理绑定)
│ └─ MessageChannelConfigurer(通道后处理器)
│
├─ 2. 发现绑定目标
│ ├─ 注解式:@EnableBinding 解析接口方法
│ └─ 函数式:Function 定义自动识别
│
├─ 3. BindingService.bindConsumer/bindProducer
│ └─ Binder 与中间件建立连接
│
└─ 4. 通道就绪,消息开始流动BinderFactory:Binder 的创建
类结构
java
// org.springframework.cloud.stream.binder.BinderFactory
public interface BinderFactory {
// 获取指定类型配置的 Binder
Binder getBinder(String configurationName, Class<? extends Binder> binderType);
// 所有已创建的 Binder
Map<String, Binder> getBinders();
}java
// 默认实现 DefaultBinderFactory
public class DefaultBinderFactory implements BinderFactory {
private final List<BinderType> binderTypes; // 已注册的 Binder 类型
private final Map<String, Binder> binderCache; // 实例缓存
@Override
public Binder getBinder(String configurationName, Class<? extends Binder> binderType) {
// 1. 根据名称查找 Binder 类型(kafka/rabbit)
BinderType binderType = binderTypes.stream()
.filter(b -> b.getName().equals(configurationName))
.findFirst()
.orElseThrow(...);
// 2. 从缓存取(每个 Binder 单例)
return binderCache.computeIfAbsent(configurationName, name -> {
// 3. 通过配置类创建 Binder 实例
return createBinder(binderType);
});
}
}Binder 的选择
yaml
spring:
cloud:
stream:
binders:
defaultBinder: # Binder 名称
type: kafka # kafka / rabbit
environment: ...classpath 有 spring-cloud-stream-binder-kafka → 注册 KafkaMessageChannelBinder
classpath 有 spring-cloud-stream-binder-rabbit → 注册 RabbitMessageChannelBinderBindingService:绑定管理
类结构
java
// org.springframework.cloud.stream.binding.BindingService
public class BindingService {
private final BinderFactory binderFactory; // Binder 工厂
private final Map<String, Binding> bindingMap; // 绑定缓存(bindingName → Binding)
private final MessageChannelConfigurer messageChannelConfigurer;
// 绑定消费者
public Binding bindConsumer(String name, String group, MessageChannel inboundTarget,
ConsumerProperties consumerProperties) {
// 1. 取对应 Binder
Binder binder = binderFactory.getBinder(properties, binderType);
// 2. 调用 Binder.bindConsumer
Binding binding = binder.bindConsumer(name, group, inboundTarget, consumerProperties);
// 3. 缓存
bindingMap.put(binding.getName(), binding);
return binding;
}
// 绑定生产者
public Binding bindProducer(String name, MessageChannel outboundTarget,
ProducerProperties producerProperties) { ... }
}绑定名称
bindConsumer: <destination>.<group> (消费绑定)
bindProducer: <destination> (生产绑定)绑定流程:消费者
BindingService.bindConsumer(name, group, inputChannel, props)
│
▼
KafkaMessageChannelBinder.bindConsumer(...)
│
├─ 1. 创建 MessageProducer(Spring Integration 组件)
│ 连接 Kafka Consumer 与应用 Input 通道
│
├─ 2. 配置:
│ ├─ topic:destination
│ ├─ group:group
│ ├─ ack 模式
│ └─ 分区分配
│
├─ 3. 启动 Kafka Consumer(后台线程拉取)
│
└─ 4. 消息到达 → MessageProducer 转为 Spring Message
→ 交给 Input 通道 → 业务方法java
// KafkaMessageChannelBinder 核心(简化)
public class KafkaMessageChannelBinder extends AbstractMessageChannelBinder<...> {
@Override
protected MessageProducer createConsumerEndpoint(...) {
// 构建 Kafka MessageProducer
KafkaInboundChannelAdapter adapter = new KafkaInboundChannelAdapter(...);
// 错误处理(重试/死信)
adapter.setErrorHandler(...);
// 关联输出通道
adapter.setOutputChannel(bindingTarget);
return adapter;
}
}绑定流程:生产者
BindingService.bindProducer(name, outputChannel, props)
│
▼
KafkaMessageChannelBinder.bindProducer(...)
│
├─ 1. 创建 MessageHandler(Spring Integration 组件)
│ 连接应用 Output 通道与 Kafka Producer
│
├─ 2. 配置:
│ ├─ topic:destination
│ ├─ 分区键表达式
│ ├─ 重试
│ └─ 确认模式
│
├─ 3. 启动 Kafka Producer
│
└─ 4. 业务发送 → Output 通道 → MessageHandler
→ 分区计算 → Kafka Producer 发送消息通道创建:DirectChannel
通道类型
java
// 默认 Input/Output 通道是 DirectChannel
public class DirectChannel extends AbstractSubscribableChannel {
// 同步分发:一个订阅者直接调用,无队列
}生产流程:
业务方法 → channel.send(message) → MessageHandler → Binder → MQ
消费流程:
MQ → MessageProducer → channel.send(message) → 业务方法(订阅者)消息流
MessageHandler(发送端)
│
├─ 分区计算(partition-key-expression → partition)
├─ 消息转换(对象 → 字节)
└─ 发送到 MQ消息转换链
序列化(生产)
业务对象 → MessageConverter.toMessage()
│
├─ 默认:JavaSerializationMessageConverter(改 contentType 用 JSON)
│ 或 CompositeMessageConverter(按 contentType 选择)
│
└─ 字节数组 → MQyaml
spring:
cloud:
stream:
bindings:
out-0:
content-type: application/json # JSON 序列化反序列化(消费)
MQ 字节 → MessageConverter.fromMessage()
│
└─ 按 contentType 还原为业务对象(方法参数类型)转换器注册
java
// org.springframework.cloud.stream.converter.CompositeMessageConverter
public class CompositeMessageConverter implements MessageConverter {
private final List<MessageConverter> converters; // 多个转换器链
@Override
public Object fromMessage(Message<?> message, Class<?> targetClass) {
// 按 contentType 选择匹配的转换器
for (MessageConverter converter : converters) {
if (converter.supports(message, targetClass)) {
return converter.fromMessage(message, targetClass);
}
}
throw new MessageConversionException(...);
}
}分区实现
生产端分区
java
// AbstractMessageChannelBinder 中
public PartitionHandler createPartitionHandler(String topic, ProducerProperties properties) {
// 根据分区键表达式计算分区号
return new PartitionHandler(expressionEvaluator, expression, partitionCount);
}partition = hash(分区键) % partitionCount消费端分区
java
// Kafka 消费端设置 partition 分配
// 分区数一致 + key 一致 → 同一分区有序消费者组与实例数
组概念源码
java
// Kafka Binder 中:
// group 对应 Kafka Consumer Group
// 同组多实例 → Kafka 自动负载均衡(分区分配给不同实例)
// 不同组 → 各自消费全部消息并发消费
yaml
spring:
cloud:
stream:
bindings:
in-0:
consumer:
concurrency: 3 # 并发消费者数(Kafka 分区数以内)完整源码链路图
生产:
业务代码 send
→ Output 通道(DirectChannel)
→ MessageHandler(分区计算 + 转换)
→ Kafka Producer(KafkaTemplate)
→ Topic
消费:
Topic
→ Kafka Consumer(分区拉取)
→ KafkaInboundChannelAdapter(MessageProducer)
→ Input 通道(DirectChannel)
→ 业务方法(@StreamListener / Consumer 函数)常见问题
- Binder 找不到? 检查 classpath 是否引入对应 binder 依赖、binders 配置 name 是否匹配。
- 消息序列化异常? contentType 与消息实际格式不匹配;对象需有默认构造与 getter。
- 分区消费乱序? 分区键表达式与分区数需两端一致,且并发数不超过分区数。
- 消费者组不生效? group 需配置在
bindings.<name>.group,否则使用匿名组(重启丢 offset)。