八
消息驱动:Spring Cloud Stream 与 RocketMQ
Message-Driven Microservices
同步调用(Feign)有个问题:下单服务要同时调用户、库存、优惠券三个服务,任何一个慢都拖慢下单接口。而且这三个服务挂了,下单就失败。消息队列的价值就是解耦、异步、削峰——下单后发一条"订单已创建"的消息出去,谁感兴趣谁自己来消费,下单接口不等。
论什么时候该用消息队列
解耦:下单后要发短信、加积分、推 BI 报表。下单服务不用知道这些下游的存在,发个消息完事,以后新增需求只加消费者。
异步:用户注册后要发欢迎邮件,邮件服务慢?没关系,发完消息立刻返回,邮件后台慢慢发。
削峰:秒杀瞬间 10 万请求打进来,数据库扛不住。消息队列当缓冲,消费者按自己的能力慢慢拉。
国内主流:RocketMQ(阿里开源),金融级可靠;Kafka 适合大数据/日志场景。Spring Cloud Stream 是一个抽象层,屏蔽底层 MQ 差异——换 RocketMQ 或 Kafka 只改 Binder,业务代码不动。
Spring Cloud Stream 概念速览(RocketMQ Binder)
# Binder:跟中间件对接的"驱动",RocketMQ / Kafka / RabbitMQ 各有一个
# Source:消息生产者,往 Topic 发消息
# Sink:消息消费者,从 Topic 收消息
# Processor:既能发又能收
# 业务代码里只要绑定一个函数,Spring Cloud Stream 自动把它接到 MQ
import java.util.function.Consumer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class OrderConsumer {
// 这个 Bean 名字 = 消费组名,自动订阅 order-topic
@Bean
public Consumer<String> orderConsumer() {
return msg -> {
System.out.println("收到订单消息:" + msg);
// 这里干发短信、加积分之类的活
};
}
}
RocketMQ 与 Kafka 怎么选
| 中间件 | 特点与适用 |
|---|---|
| RocketMQ | 阿里开源,金融级可靠,支持事务消息、延迟消息、死信队列,国内电商/金融主流。中文文档全。 |
| Kafka | 高吞吐,日志/大数据场景王者(埋点、流处理)。业务消息也能做,但事务消息能力弱于 RocketMQ。 |
| RabbitMQ | 老牌,路由灵活,吞吐量不如前两个,中小项目用。 |
消息幂等:同一条消息不能扣两次钱
MQ 可能因为网络重试,把同一条消息投递给你两次。如果消费逻辑是"给用户加 10 块钱",收到两次就加了 20 块——事故。幂等的意思是:同一条消息处理一次和处理十次,结果一样。常见做法:给每条消息一个唯一 ID,消费前先查这个 ID 处理过没有,处理过就跳过。
幂等消费的思路(伪代码)
@Bean
public Consumer<OrderMsg> orderConsumer() {
return msg -> {
String msgId = msg.getMsgId();
// 1. 先查这个 msgId 处理过没有
if (processedDao.exists(msgId)) {
System.out.println("重复消息,跳过:" + msgId);
return;
}
// 2. 正常业务逻辑
sendSms(msg.getPhone());
// 3. 标记已处理
processedDao.save(msgId);
};
}
消息可靠性与死信队列
消息发到 MQ 不等于业务处理成功。消费者处理失败(比如发邮件接口挂了),不能直接把消息丢了——要重试。重试 N 次还失败,就把消息丢进死信队列(DLQ),人工介入。生产环境必做:① 生产端确认消息真的到 MQ;② 消费端手动 ACK;③ 失败重试 + 死信兜底;④ 消息幂等(同一条消息消费两次结果一样)。