Skip to content

消息队列

消息队列(Message Queue,简称 MQ)是分布式系统中非常常见的基础组件。它的核心作用不是“让系统变复杂”,而是把原来必须同步完成的一段调用,拆成“先把消息可靠写进去,再由下游慢慢处理”的异步链路。

通俗理解:系统 A 不直接调用系统 B,而是把一张“任务单”放到队列里;系统 B 有空时从队列取任务并处理。这样可以降低系统之间的耦合,也能在流量突然变大时保护下游服务。

为什么需要消息队列

1. 异步解耦

在没有 MQ 的情况下,订单服务创建订单后可能要同步调用库存、积分、短信、物流等服务。任何一个服务慢或者失败,都会影响创建订单接口。

使用 MQ 后,订单服务只需要把“订单已创建”事件发送出去,下游系统各自订阅并处理,主流程更短,也更容易扩展。

mermaid
flowchart TD
    A["订单服务"] --> B["消息队列"]
    B --> C["库存服务"]
    B --> D["积分服务"]
    B --> E["短信服务"]
    B --> F["物流服务"]

2. 削峰填谷

秒杀、活动、支付回调这类场景会出现瞬时高峰。如果请求全部直接打到数据库或下游服务,很容易把系统压垮。MQ 可以先承接高峰流量,下游按自己的处理能力稳定消费。

mermaid
flowchart TD
    A[突发请求] --> B[写入 MQ]
    B --> C[消费者按固定速率处理]
    C --> D[数据库压力平稳]

3. 最终一致性

分布式系统里很难让多个服务同时成功或同时失败。MQ 更常见的做法是保证主业务先成功,再通过可靠消息、重试、补偿等机制让下游最终处理成功。

核心概念

概念说明
Producer生产者,负责发送消息
Consumer消费者,负责接收并处理消息
Broker消息服务器,负责接收、存储、投递消息
Topic / Queue消息分类或队列,不同产品命名略有差异
Offset消费进度,记录消费者处理到哪里
ACK消费确认,表示消息已经被成功处理
Retry重试,消费失败后再次投递
DLQ死信队列,多次失败后进入的异常队列

消息队列的整体流程

mermaid
flowchart TD
    P["生产者构建消息"] --> S["发送到 Broker"]
    S --> W["Broker 持久化"]
    W --> R["返回发送结果"]
    W --> Q["等待消费者拉取或推送"]
    Q --> C["消费者处理业务"]
    C --> A{"处理成功?"}
    A -->|是| ACK["提交消费确认"]
    A -->|否| Retry["进入重试或死信"]

常见问题

消息会不会丢

消息是否会丢,取决于三个阶段:

  • 生产阶段:生产者发送失败是否重试,是否等待 Broker 确认。
  • 存储阶段:Broker 是否刷盘,是否有主从复制。
  • 消费阶段:消费者是否在业务真正成功后再确认。

所以“用 MQ 就不会丢消息”是不准确的。正确做法是根据业务选择可靠发送、持久化、手动 ACK、消费幂等等策略。

消息会不会重复

多数 MQ 更容易保证“至少投递一次”,也就是说消息可能重复。网络超时、消费者处理成功但 ACK 失败、Broker 重试等情况都会造成重复投递。

因此消费端必须设计幂等。常见方式包括唯一业务号去重、数据库唯一索引、状态机判断、Redis setnx 等。

消息顺序如何保证

顺序消息通常只能在一个局部范围内保证,例如同一个订单的消息进入同一个队列,由同一个消费者顺序处理。全局顺序会牺牲吞吐量,实际业务中很少使用。

消息堆积怎么办

消息堆积不是简单“加消费者”就能解决。要先判断生产速度、消费速度、消费者错误率、重试量、下游数据库或接口耗时、分区/队列是否热点,再决定限流、扩容、批处理、失败隔离、死信、补偿还是调整分区队列。

详细学习:消息堆积与背压

目录

从零到生产级掌握

MQ 主线课程。把生产者、Broker、消费者、ACK、Offset、可靠性、重复消费、幂等、顺序消息、事务消息、重试、死信、堆积、背压、RabbitMQ/Kafka/RocketMQ 选型串成完整体系。

从零到精通验收清单

用可验证任务判断自己是不是真的学会 MQ。它按作用、发送、Broker、消费、可靠性、幂等、顺序、事务、死信、堆积、Lag、背压、三大 MQ 选型和商业场景拆成验收标准。

商业场景训练营

用订单创建事件、Outbox 本地消息表、消费者幂等、业务成功后 ACK、ES 同步失败补偿、消息堆积、扩容反弹、Lag 分布、RabbitMQ Ready/Unacked 和三大 MQ 选型,把原理放到商业项目里跑起来。

消息堆积与背压

生产排障高频知识点。重点学习消息积压、消费 Lag、Ready/Unacked、重试堆积、背压、限流、扩容、死信和补偿。

扩容后再次堆积

专门解释“消费者扩容后短暂有效,后面又开始堆积”的原理。重点学习有效消费 TPS、分区/队列并行度、下游瓶颈、重试风暴、热点 key、线程池和连接池排查。

Lag 分布与堆积定位

专门解释“Lag 分布是什么”。重点学习为什么不能只看总 Lag,而要按 Kafka Partition、RocketMQ MessageQueue、RabbitMQ Queue/Unacked 分布判断整体消费慢、局部热点、慢消息、重试风暴和顺序阻塞。

MQ面试题

刷题复习入口。详细原理回到对应知识点页学习。

RabbitMQ

适合传统企业应用、复杂路由、任务分发、延迟/死信处理、对 AMQP 生态有要求的场景。

Kafka

适合高吞吐事件流、日志采集、用户行为埋点、CDC 数据同步、实时计算、数据管道和可回放场景。

RocketMQ

适合高吞吐、金融级可靠消息、事务消息、顺序消息、消费重试等场景。

官网:为什么选择 RocketMQ | RocketMQ

RabbitMQ、Kafka 和 RocketMQ 怎么选

对比项RabbitMQKafkaRocketMQ
核心模型Exchange + Queue + BindingTopic + Partition + LogTopic + Queue + ConsumerGroup
吞吐能力中高吞吐,路由能力强高吞吐,适合连续事件流高吞吐,适合业务消息堆积
消息保留偏队列模型,消费后可删除按时间或大小保留,天然支持回放支持堆积、重试和死信
顺序消息单队列可做局部顺序分区内顺序队列内顺序,业务语义更直接
事务消息通常靠发布确认、业务表、补偿实现支持 Kafka 内部事务,外部库需 Outbox/CDC原生事务消息能力较强
路由能力很强,Exchange 类型丰富较弱,主要靠 Topic 和 key中等,靠 Topic、Tag、Key
适合场景企业集成、任务队列、复杂路由、延迟/死信日志、埋点、数据同步、实时计算、事件总线电商交易、金融异步链路、事务消息、顺序消息

学习建议:先理解本页的 MQ 通用概念,再学习 RabbitMQ 核心模型Kafka 架构与存储原理RocketMQ 领域模型。如果重点关注分布式事务,直接看 RocketMQ 事务消息;如果重点关注日志、数据同步和回放,直接看 Kafka 商业常用场景

代码 Demo:订单创建后发送消息

下面用伪代码演示 MQ 在业务里的位置。重点是:主业务先落库,再发送事件消息,下游异步处理。

java
@Service
public class OrderService {
    private final OrderRepository orderRepository;
    private final MessageProducer messageProducer;

    public OrderService(OrderRepository orderRepository, MessageProducer messageProducer) {
        this.orderRepository = orderRepository;
        this.messageProducer = messageProducer;
    }

    @Transactional
    public Long createOrder(CreateOrderCommand command) {
        Order order = new Order(command.userId(), command.productId(), command.count());
        orderRepository.save(order);

        OrderCreatedEvent event = new OrderCreatedEvent(order.getId(), order.getUserId());
        messageProducer.send("order.created", event);
        return order.getId();
    }
}

消费端必须做幂等:

java
@Component
public class StockConsumer {
    public void onMessage(OrderCreatedEvent event) {
        if (alreadyProcessed(event.orderId())) {
            return;
        }
        reduceStock(event.orderId());
        markProcessed(event.orderId());
    }
}

真实项目可以用 RabbitMQ、Kafka、RocketMQ 等产品实现 MessageProducer,但“可靠发送、业务幂等、失败重试”这三个原则都不能省。