消息队列基础
定义与作用
消息队列(Message Queue, MQ)是一种异步通信中间件,用于在分布式系统的不同组件之间传递消息。它的核心价值在于解耦:生产者(Producer)发送消息后无需等待消费者(Consumer)立即处理,消费者按自己的节奏消费消息。
在分布式系统中,消息队列解决了三个核心矛盾:
| 矛盾 | 没有消息队列 | 引入消息队列后 |
|---|---|---|
| 紧耦合 | 服务 A 直接调用服务 B,B 宕机则 A 失败 | A 只与 MQ 交互,B 恢复后从 MQ 消费 |
| 流量冲击 | 瞬时高峰直接打到下游,触发雪崩 | MQ 缓冲消息,消费者匀速处理(削峰填谷) |
| 同步延迟 | A 必须等 B 完成才能返回 | A 发送后立即返回,异步处理 |
两种经典模型
点对点(Point-to-Point / Queue)
生产者向队列发送消息,有且仅有一个消费者消费该消息。消息被消费后从队列中移除。
发布-订阅(Publish-Subscribe / Topic)
生产者向 Topic 发布消息,所有订阅该 Topic 的消费者都收到一份消息副本。
Kafka 的模型融合了两种范式:Topic 层面对应发布-订阅,消费组(Consumer Group)内部对应点对点(同一分区只被组内一个消费者消费)。这一设计是 Kafka 能够同时支撑"广播"和"负载均衡"两种场景的关键。
核心场景
场景一:日志收集(解耦 + 高吞吐)
互联网公司每天产生 TB 级别的用户行为日志。前端埋点 SDK → Kafka → 多个下游(Hadoop 离线分析、Flink 实时计算、ES 搜索索引)。如果没有 Kafka,每个下游都需要独立对接日志源。
场景二:秒杀系统(削峰填谷)
双十一零点瞬间涌入百万并发请求。若直接透传到订单系统,数据库瞬间被打满。通过前端限流 + Kafka 缓冲,订单服务以自身处理能力(如 5000 QPS)匀速消费,多余请求排队等待或快速失败。
场景三:微服务异步解耦
用户注册事件:UserService → Kafka → EmailService(发欢迎邮件)、CouponService(发新人优惠券)、RecommendationService(初始化推荐模型)。UserService 不需要知道有哪些下游。
简单实现与消息队列的对比
操作前(同步直接调用):
// UserService 直接耦合所有下游
public void register(String username) {
userDao.insert(username);
// 必须等邮件发送完成才返回 —— 耦合且慢
emailService.sendWelcomeEmail(username);
// 若此处抛异常,用户已经注册但优惠券没发
couponService.sendNewUserCoupon(username);
}
操作后(引入消息队列):
// UserService 只与 Kafka 交互
public void register(String username) {
userDao.insert(username);
producer.send(new UserRegisterEvent(username));
// 立即返回,下游异步消费
}
状态对比:
| 维度 | 同步直接调用 | 引入消息队列 |
|---|---|---|
| 响应延迟 | 300ms+(等所有下游) | < 5ms(仅写 Kafka) |
| 下游故障影响 | 注册接口不可用 | 注册正常,故障恢复后继续消费 |
| 新增下游成本 | 修改 UserService 代码 | 新增一个 Consumer 订阅即可 |
| 数据一致性 | 强一致(但耦合) | 最终一致(需额外设计补偿) |
易错场景
场景:将消息队列当作数据库使用,期望消息持久化到永远。
问题:Kafka 虽然有持久化特性(与 RabbitMQ 等不同),但消息默认有保留期限(log.retention.hours=168,7天)。超期自动删除。不要把生命周期超过保留期的数据只存放在 Kafka 中。
面试高频考点
Q:消息队列是银弹吗?引入 MQ 带来哪些新问题?
A:不是。引入 MQ 带来:
- 系统可用性降低 — MQ 本身成为关键路径,需要考虑 MQ 的高可用(集群 + 副本)
- 复杂度上升 — 需要处理消息重复、消息丢失、消息顺序等一致性问题
- 最终一致性 — 放弃强一致,需要设计补偿机制和幂等处理