消息中间件与 Spring 集成深入
Kafka 与 RabbitMQ 是两种主流消息中间件,Spring Cloud Stream 的 Binder 把它们接入统一模型。本文拆解两个 Binder 的实现差异,并给出可靠性投递与幂等消费的完整方案。
消息中间件对比
| 维度 | Kafka | RabbitMQ |
|---|---|---|
| 模型 | Topic + Partition(日志流) | Exchange + Queue(路由) |
| 顺序保证 | 分区内有序 | 单队列有序 |
| 吞吐 | 极高(顺序写磁盘) | 较高 |
| 消费模型 | 拉取(offset 管理) | 推送(ACK 确认) |
| 延迟 | 毫秒级 | 微秒级 |
| 消息堆积 | 天然支持(磁盘日志) | 需队列配置 |
| 重试 | 手动(seek offset) | DLX/重试队列 |
Kafka: 生产者 → Topic(分区) → 消费者组(拉取+offset)
Rabbit: 生产者 → Exchange → Queue → 消费者(推送+ACK)Kafka Binder 源码分析
核心类
java
// org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder
public class KafkaMessageChannelBinder extends AbstractMessageChannelBinder<
ConsumerProperties, ProducerProperties, KafkaBindingProperties> {
...
}消费端
java
// 构建 Kafka Consumer
@Override
protected MessageProducer createConsumerEndpoint(...) {
// 1. 构建 Kafka Consumer 配置
KafkaConsumerProperties consumerProperties = ...
// 2. 创建 KafkaInboundChannelAdapter
KafkaInboundChannelAdapter adapter = new KafkaInboundChannelAdapter(
kafkaMessageListenerContainer, listenerProperties);
// 3. 配置 ACK 模式(自动/手动)
adapter.setAckMode(consumerProperties.getAckMode());
// 4. 错误处理(重试/死信)
adapter.setErrorHandler(new SeekToCurrentErrorHandler(
new DeadLetterPublishingRecoverer(kafkaTemplate), retries));
return adapter;
}Kafka Consumer 配置
yaml
spring:
cloud:
stream:
kafka:
binder:
brokers: localhost:9092
auto-create-topics: true
replication-factor: 1
bindings:
in-0:
consumer:
ack-mode: manual # 手动 ACK
auto-commit-offset: false
start-offset: latest # earliest/latest
enable-dlq: true # 死信队列消息确认机制
Kafka 拉取消息 → 业务处理 → commit offset(提交消费位置)
│
├─ 处理成功 → 提交 offset(下一条继续)
├─ 处理失败(可重试)→ 重试 → 成功提交
└─ 处理失败(不可恢复)→ 死信队列(-dlt 后缀)死信队列
原 Topic:order.topic
死信 Topic:order.topic.order-group.DLT
触发条件:超过最大重试次数(max-attempts)
死信消息头保留原始信息(原 topic、partition、offset)RabbitMQ Binder 源码分析
核心类
java
// org.springframework.cloud.stream.binder.rabbit.RabbitMessageChannelBinder
public class RabbitMessageChannelBinder extends AbstractMessageChannelBinder<...> {
...
}消费端
java
@Override
protected MessageProducer createConsumerEndpoint(...) {
// 1. 构建 RabbitMQ 监听容器
SimpleMessageListenerContainer container = createListenerContainer(...);
// 2. 创建 AmqpInboundChannelAdapter
AmqpInboundChannelAdapter adapter = new AmqpInboundChannelAdapter(container);
// 3. ACK 模式:自动确认/手动确认
adapter.setAcknowledgeMode(consumerProperties.getAcknowledgeMode());
// 4. 声明队列与绑定(Exchange → Queue)
RabbitAdmin.declareQueue(...);
return adapter;
}RabbitMQ 路由结构
Exchange(交换机)
│ routing key
▼
Queue(队列)
│
▼
消费者(ACK)Spring Cloud Stream 默认:
Exchange:<destination>(topic 类型)
Queue: <destination>.<group>
绑定: routing key = #重试与死信
yaml
spring:
cloud:
stream:
rabbit:
bindings:
in-0:
consumer:
acknowledge-mode: manual # 手动确认
max-attempts: 3 # 重试次数
dlq-enabled: true # 死信队列RabbitMQ 死信:
Queue 设置 x-dead-letter-exchange 属性
→ 失败消息进入 DLX → DLQ(.dlq 后缀)手动 ACK
java
@StreamListener("in-0")
public void onMessage(Message msg) {
try {
process(msg);
// 手动确认
((Channel) msg.getHeaders().get(AmqpHeaders.CHANNEL))
.basicAck(msg.getHeaders().get(AmqpHeaders.DELIVERY_TAG, Long.class), false);
} catch (Exception e) {
// 拒绝(重新入队或进死信)
((Channel) msg.getHeaders().get(AmqpHeaders.CHANNEL))
.basicReject(deliveryTag, false);
}
}Kafka vs Rabbit Binder 对比
| 维度 | Kafka Binder | Rabbit Binder |
|---|---|---|
| 消息模型 | Topic/Partition | Exchange/Queue |
| ACK | offset 提交 | basicAck/basicNack |
| 重试 | max-attempts + 手动 seek | 重试队列 + DLX |
| 死信 | Topic.DLT | DLX → DLQ |
| 顺序 | 分区内 | 单队列 |
| 积压处理 | 天然好 | 需监控 |
消息可靠性投递
可靠性全景
生产者可靠性 ──▶ Broker 可靠性 ──▶ 消费者可靠性
│ │ │
├─ 发送确认 ├─ 持久化 ├─ ACK 机制
├─ 重试 ├─ 副本同步 ├─ 幂等消费
└─ 事务 └─ 容错恢复 └─ 死信处理生产端可靠性
yaml
spring:
kafka:
producer:
acks: all # 全副本确认
retries: 3 # 重试
enable-idempotence: true # 幂等生产者(防重复消息)java
// 发送确认回调
kafkaTemplate.send(topic, data).addCallback(
success -> log.info("发送成功: {}", success.getRecordMetadata().offset()),
failure -> {
// 发送失败处理(重试/记录)
log.error("发送失败", failure);
});Broker 端可靠性
Kafka:replication-factor(副本数)
min.insync.replicas(最小同步副本)
Rabbit:持久化 Exchange/Queue(durable)
持久化消息(delivery_mode=2)消费端可靠性
yaml
spring:
kafka:
consumer:
enable-auto-commit: false # 手动提交java
// 手动 ACK(Kafka)
@KafkaListener(topics = "order.topic")
public void onMessage(ConsumerRecord<String, String> record,
Acknowledgment ack) {
try {
process(record.value());
ack.acknowledge(); // 处理成功才提交 offset
} catch (Exception e) {
// 记录失败,可重试或进死信
}
}幂等消费方案
为什么必须幂等
消息重发场景:
1. 生产者重试 → 同一条消息发多次
2. 消费者处理成功但 ACK 前崩溃 → offset 未提交 → 重启重发
3. 网络抖动 → 重复投递
结果:业务被重复执行(重复扣款、重复下单)方案一:唯一键约束
java
// 业务表唯一键 = 消息 ID(幂等键)
@Insert("INSERT INTO order_oper_log (id, order_id, op) VALUES (#{msgId}, ...)")
// 插入失败(主键冲突)→ 说明已处理,直接跳过java
public void onMessage(OrderMessage msg) {
try {
orderService.apply(msg); // 幂等:内部用 msg.id 做唯一键
} catch (DuplicateKeyException e) {
log.info("重复消息,跳过: {}", msg.getId());
}
}方案二:Redis 去重
java
public void onMessage(OrderMessage msg) {
String key = "msg:processed:" + msg.getId();
// SETNX:已存在说明处理过
Boolean first = redisTemplate.opsForValue()
.setIfAbsent(key, "1", Duration.ofHours(24));
if (!first) {
log.info("重复消息,跳过");
return;
}
process(msg); // 首次处理
}方案三:状态机幂等
java
// 订单状态机:只有"待支付"能转到"已支付"
public void onPayment(PaymentMessage msg) {
boolean ok = orderService.updateStatus(
msg.getOrderId(),
OrderStatus.PENDING_PAY, // 期望当前状态
OrderStatus.PAID); // 目标状态
if (!ok) {
// 状态已变更(说明已处理过),跳过
log.info("订单状态已变更,跳过");
}
}幂等方案对比
| 方案 | 实现难度 | 适用范围 |
|---|---|---|
| 唯一键 | 低 | 有数据库写操作 |
| Redis 去重 | 低 | 无状态服务 |
| 状态机 | 中 | 有状态流转业务 |
| 消息去重表 | 中 | 通用场景 |
生产兜底方案
消息对账
定时任务:扫描本地消息表/订单表
├─ 未发送/未处理的消息 → 重新投递
└─ 对账失败 → 告警人工介入监控指标
Producer:发送成功率、重试次数、积压量
Broker: topic 积压、分区倾斜
Consumer:消费速率、失败率、DLT 数量常见问题
- 消息重复但业务无法幂等? 用唯一键/去重表兜底,任何消息系统都可能重复。
- 积压怎么办? 扩容消费者(分区数内);临时降级消费逻辑;查询链路分析。
- 手动 ACK 忘提交? 处理成功必须 ack,否则消息永远不被消费(Kafka 会重新投递)。
- 死信堆积? 死信消费任务定期处理;结合对账任务重放或人工介入。