消息队列在现代分布式系统中承担着削峰填谷、异步解耦、事件驱动等重要职责。Spring 生态对 RabbitMQ 和 Kafka 都提供了深度的整合支持,开发者无需手动管理连接、通道或消费者线程,即可快速集成消息能力。然而,引入消息队列后,消息的可靠投递就成为必须严肃对待的课题——消息丢失、重复消费、顺序错乱等问题一旦出现,排查和修复成本极高。
本节将分别介绍 Spring Boot 如何整合 RabbitMQ 和 Kafka,涵盖基本收发、序列化配置以及可靠投递的关键实践。
14.4.1 RabbitMQ 基础整合
1. 依赖与配置
在 Spring Boot 项目中引入 RabbitMQ 支持,只需添加起步依赖:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
然后在 application.yml 中配置连接信息:
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
virtual-host: /
# 以下为生产环境建议配置
publisher-confirm-type: correlated # 开启发送方确认
publisher-returns: true # 开启路由失败回退
template:
mandatory: true # 消息无法路由时回调
listener:
simple:
acknowledge-mode: manual # 手动确认模式(默认 auto)
retry:
enabled: true
max-attempts: 3
initial-interval: 1000ms
2. 声明队列、交换器与绑定
推荐通过配置类显式声明消息系统的拓扑结构,避免依赖消费者自动创建带来的不可控风险:
@Configuration
public class RabbitConfig {
public static final String ORDER_EXCHANGE = "order.exchange";
public static final String ORDER_QUEUE = "order.queue";
public static final String ORDER_ROUTING_KEY = "order.created";
@Bean
public DirectExchange orderExchange() {
return new DirectExchange(ORDER_EXCHANGE, true, false);
}
@Bean
public Queue orderQueue() {
return QueueBuilder.durable(ORDER_QUEUE).build();
}
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue())
.to(orderExchange())
.with(ORDER_ROUTING_KEY);
}
}
3. 发送消息
使用 Spring 提供的 RabbitTemplate 即可完成消息发送:
@Component
public class OrderMessageSender {
private final RabbitTemplate rabbitTemplate;
public OrderMessageSender(RabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;
}
public void sendOrderCreated(Order order) {
// 可将对象转为 JSON 后发送,需要配置 MessageConverter
rabbitTemplate.convertAndSend(
RabbitConfig.ORDER_EXCHANGE,
RabbitConfig.ORDER_ROUTING_KEY,
order
);
}
}
为了让对象自动序列化为 JSON,需配置 MessageConverter:
@Bean
public MessageConverter jsonMessageConverter() {
return new Jackson2JsonMessageConverter();
}
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplate template = new RabbitTemplate(connectionFactory);
template.setMessageConverter(jsonMessageConverter());
return template;
}
4. 接收消息
消费者通过 @RabbitListener 注解绑定队列,消息体的反序列化由消息转换器自动完成:
@Component
public class OrderMessageListener {
@RabbitListener(queues = RabbitConfig.ORDER_QUEUE)
public void handleOrderCreated(Order order, Channel channel, Message message)
throws IOException {
try {
// 执行业务逻辑...
System.out.println("收到订单创建消息:" + order.getId());
// 手动确认
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
} catch (Exception e) {
// 根据业务决定是否重新入队
channel.basicNack(message.getMessageProperties().getDeliveryTag(),
false, true); // 重新入队
}
}
}
在 acknowledge-mode: manual 下,必须手动调用 basicAck 确认消息,否则消息会一直处于未确认状态,直到连接断开后重新入队。
14.4.2 RabbitMQ 可靠投递实践
1. 发送方确认(Publisher Confirm)
仅将消息写入 Socket 缓冲区并不代表 Broker 已成功接收。开启 publisher-confirm-type: correlated 后,可为每条消息设置确认回调:
@Component
public class ReliableOrderSender {
private final RabbitTemplate rabbitTemplate;
public ReliableOrderSender(RabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;
}
public void sendOrderWithConfirm(Order order) {
CorrelationData correlationData = new CorrelationData(order.getId());
// 设置确认回调
correlationData.getFuture().whenComplete((confirm, throwable) -> {
if (throwable != null) {
// 发送异常,记录日志,启动补偿
log.error("消息发送异常: {}", order.getId(), throwable);
compensationService.retryOrAlert(order);
} else if (confirm.isAck()) {
log.info("消息已确认到达Broker: {}", order.getId());
} else {
log.warn("消息被Broker拒绝: {}", order.getId());
// 可实现重发或告警
}
});
rabbitTemplate.convertAndSend(
RabbitConfig.ORDER_EXCHANGE,
RabbitConfig.ORDER_ROUTING_KEY,
order,
correlationData
);
}
}
2. 路由失败处理
若消息发布到交换器后无法路由到任何队列(例如路由键写错),开启 mandatory: true 并注册 ReturnCallback 可捕获这类失败:
@PostConstruct
public void initReturnCallback() {
rabbitTemplate.setReturnsCallback(returned -> {
log.warn("消息无法路由: 交换器={} 路由键={} 消息内容={}",
returned.getExchange(), returned.getRoutingKey(),
new String(returned.getMessage().getBody()));
// 可存入数据库或重发
});
}
3. 消费者手动确认与重试
消费者侧最常见的问题是消费失败后直接丢失消息,或陷入无限重试的死循环。通过手动确认模式结合重试,可以实现精细化控制:
- 业务异常(如库存不足)不应重试,应当在代码中 catch 后直接确认并记录或发送补偿事件。
- 临时故障(如网络超时、数据库死锁)可以重试,但要限制次数并设置退避策略。
@RabbitListener(queues = RabbitConfig.ORDER_QUEUE)
public void handleWithRetry(Order order, Channel channel, Message message) throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
processOrder(order);
channel.basicAck(deliveryTag, false);
} catch (RetryableException e) {
// 临时故障,重试有限次
if (hasExceedMaxRetry(message)) {
channel.basicNack(deliveryTag, false, false); // 拒绝且不重新入队
compensationService.handleFailure(order, e);
} else {
channel.basicNack(deliveryTag, false, true); // 重新入队
}
} catch (Exception e) {
// 非重试异常,直接拒绝并记录
channel.basicNack(deliveryTag, false, false);
compensationService.handleFailure(order, e);
}
}
可以通过消息头记录重试次数,每次重新入队时次数递增。
4. 消息幂等性保证
即使消息确认机制再完善,网络抖动仍可能导致 Broker 重复投递。消费者必须实现幂等处理:
- 为每条消息分配业务唯一标识(如订单号),并在消费前检查 Redis 或数据库是否已处理。
- 利用数据库唯一约束,将“处理消息”与“记录处理状态”放在同一事务中。
14.4.3 Kafka 基础整合
1. 依赖与配置
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all # 等待所有副本确认
retries: 3 # 发送重试
enable-idempotence: true # 幂等发送,配合 acks=all
consumer:
group-id: order-group
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "com.example.dto"
enable-auto-commit: false # 关闭自动提交,手动控制位移
listener:
ack-mode: manual # 手动确认
2. 发送消息
使用 KafkaTemplate 发送消息:
@Component
public class OrderKafkaSender {
private final KafkaTemplate<String, Order> kafkaTemplate;
public OrderKafkaSender(KafkaTemplate<String, Order> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void sendOrder(String key, Order order) {
ListenableFuture<SendResult<String, Order>> future =
kafkaTemplate.send("order-topic", key, order);
future.addCallback(
result -> log.info("消息发送成功: offset={}",
result.getRecordMetadata().offset()),
ex -> {
log.error("消息发送失败", ex);
// 启动补偿
}
);
}
}
3. 接收消息
@Component
public class OrderKafkaListener {
@KafkaListener(topics = "order-topic")
public void onMessage(ConsumerRecord<String, Order> record,
Acknowledgment ack) {
try {
Order order = record.value();
processOrder(order);
ack.acknowledge(); // 手动提交位移
} catch (Exception e) {
log.error("消费失败: offset={}", record.offset(), e);
// 不提交位移,消息会再次消费;注意避免无限重试
}
}
}
14.4.4 Kafka 可靠投递与消费处理
1. 生产者可靠性
acks=all:确保所有 ISR 副本确认后才认为消息发送成功,避免 Leader 宕机丢失。enable-idempotence=true:配合acks=all,Broker 会对生产者消息去重,防止重试导致重复。- 发送结果异步回调:记录失败消息或重试,不能捕获异常就忽略。
- 合理设置
retries与delivery.timeout.ms:对于临时故障进行重试,但要设置超时防止长时间阻塞。
2. 消费者手动位移提交
Kafka 消费者的位移管理决定消息的投递语义。关闭自动提交后,应在处理完业务逻辑后再提交位移,实现至少一次语义。若要精确一次,需将业务结果与位移存放在同一个事务中(使用 Kafka 事务或外部协调)。
@KafkaListener(topics = "order-topic")
public void onMessage(ConsumerRecord<String, Order> record,
Acknowledgment ack) {
try {
orderService.processAndMarkProcessed(record.value(), record.offset());
ack.acknowledge();
} catch (DuplicateException e) {
// 已处理过,直接提交位移,避免重复处理
ack.acknowledge();
} catch (Exception e) {
// 记录错误,不提交位移,线程挂起或暂停一段时间后继续
throw e; // 配置中的 SeekToCurrentErrorHandler 会处理
}
}
推荐配置 SeekToCurrentErrorHandler 配合 DeadLetterPublishingRecoverer,将消费多次失败的消息转移到死信主题,避免阻塞分区:
@Bean
public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(
ConsumerFactory<Object, Object> consumerFactory,
KafkaTemplate<Object, Object> template) {
ConcurrentKafkaListenerContainerFactory<Object, Object> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory);
factory.setCommonErrorHandler(new DefaultErrorHandler(
new DeadLetterPublishingRecoverer(template),
new FixedBackOff(1000L, 3) // 重试3次,间隔1秒
));
return factory;
}
3. 幂等性
即使启用了生产者的幂等和消费者手动提交,业务层仍需考虑重复消费的情况。全局唯一标识 + 业务状态检查 是通用方案。例如在消息体中携带事件 ID,消费前在 Redis 中尝试设置一个具有过期时间的键(SETNX)或插入一张去重表。
void processOrder(Order order) {
String eventId = order.getEventId();
if (redisTemplate.opsForValue().setIfAbsent(eventId, "1", Duration.ofMinutes(30))) {
// 未处理,执行核心业务
businessLogic(order);
} else {
log.info("重复消息忽略: eventId={}", eventId);
}
}
14.4.5 实践建议
- 选择合适的消息队列:RabbitMQ 擅长复杂的路由和消息管理,适合需要投递确认、死信队列等场景;Kafka 擅长高吞吐的日志和流式处理,适合事件溯源和实时数据管道。
- 消息体设计:尽量包含业务唯一 ID、事件产生时间戳和必要上下文,便于幂等、监控和补偿。
- 监控与告警:监控消息积压、消费延迟、确认失败率等指标,及时发现问题。
- 消费端做好幂等:不要完全依赖 Broker 的投递语义,业务层的幂等是最后的防线。
- 死信与补偿:对无法成功处理的消息,导入死信队列或数据库,并提供管理界面进行人工干预或定时补偿。
消息队列的整合看似简单,但可靠投递涉及生产端确认、消费端提交、重试策略、幂等设计等多个环节。Spring 的整合让我们能聚焦于这些可靠性考量本身,而不必陷入底层 API 的细节中。在实际项目中,根据业务对一致性的要求,选择合适级别的保障方案,是架构师和开发者必须掌握的技能。