RabbitMQ 基础篇

主线

同步调用的三个问题 → 异步调用 + MQ 三个角色 → RabbitMQ 核心概念
        ↓
SpringAMQP:RabbitTemplate 发、@RabbitListener 收
        ↓
五种模型:简单队列 → WorkQueue → Fanout → Direct → Topic
        ↓
自动声明队列/交换机 + JSON 消息转换器
        ↓
业务落地:支付成功异步通知交易服务

一、为什么用 MQ

同步调用(OpenFeign 链式调用)的三个问题:

问题表现原因
拓展性差支付成功后每多一件事,支付代码就要改一次主业务被迫感知旁支业务,违反开闭原则
性能下降总耗时 = 所有远程调用耗时之和同步调用阻塞等待
级联失败交易服务或通知服务挂了,事务回滚,支付失败非核心业务连累核心业务

第三条最要命:钱都扣了,没道理因为短信发不出去就退回去。这类旁支业务不该放进核心事务。

异步调用用消息通知代替直接调用,角色变成三个:消息发送者(publisher)、消息 Broker(管理、暂存、转发消息)、消息接收者(consumer)。改造后支付服务只做扣款和改流水状态,发一条消息给 Broker,交易、通知、积分各自订阅处理。

好处是耦合低、总耗时只算核心业务、加新业务不用改支付代码、下游挂了也不影响支付。代价是完全依赖 Broker 的可靠性,链路从一条调用栈变成一条消息流,排查问题更麻烦。

选型上课堂用 RabbitMQ,国内用得最多,可用性、可靠性、吞吐、延迟都比较均衡。Kafka 吞吐最高,RocketMQ 吞吐和可靠性兼顾,追求低延迟可以选 RabbitMQ 或 Kafka。


二、RabbitMQ 核心概念

概念说明
publisher生产者,发消息的一方
consumer消费者,收消息的一方
queue队列,存储消息。消息在队列里等着被消费
exchange交换机,负责路由。生产者把消息发给交换机,由它决定投到哪个队列
virtual host虚拟主机,数据隔离单位。每个 vhost 有自己的 exchange 和 queue

一句话记住:交换机只管路由不存消息,队列只管存储不认路由,两者靠 Binding 绑定。

在控制台上做过一次验证:直接往没有绑定队列的交换机发消息,消息丢了;建好队列但没绑定,消息还是到不了;绑定之后才进队列。这正好反证了上面那句话。

两个端口别搞混:5672 是代码里收发消息的通信端口,15672 是浏览器上的管理控制台。yml 里写成 15672 会直接连接异常。

多项目共用一套集群时,给每个项目建独立账号和独立 virtual host,代码里用 virtual-host 选定自己的环境。


三、SpringAMQP

加 spring-boot-starter-amqp 依赖后,自动装配会创建好 ConnectionFactory、RabbitTemplate、AmqpAdmin 和 SimpleRabbitListenerContainerFactory,配置写好地址就能直接注入使用。

publisher 和 consumer 两端的 yml 一样:

spring:
  rabbitmq:
    host: 192.168.150.101 # 虚拟机 IP
    port: 5672            # 通信端口,不是 15672
    virtual-host: /hmall
    username: hmall
    password: 123

3.1 最简单模型:直接发到队列

publisher 端:

@SpringBootTest
public class SpringAmqpTest {
    @Autowired
    private RabbitTemplate rabbitTemplate;
 
    @Test
    public void testSimpleQueue() {
        rabbitTemplate.convertAndSend("simple.queue", "hello, spring amqp!");
    }
}

consumer 端:

@Component
public class SpringRabbitListener {
    @RabbitListener(queues = "simple.queue")
    public void listenSimpleQueueMessage(String msg) {
        System.out.println("消费者接收到消息:【" + msg + "】");
    }
}

几个必须弄清楚的点:

  • convertAndSend 里的 convert,指先用 MessageConverter(默认 SimpleMessageConverter)把 Java 对象转成 AMQP 的 Message 再发出去。
  • 只传队列名也能发出去,是因为内部用了默认交换机(名字是空字符串 "")。每个队列都会自动绑定到它,bindingKey 等于队列名,routingKey 就是队列名。所谓”直接发队列”本质还是走了交换机,只是被隐藏了。这种用法一般只在测试里出现。
  • @SpringBootTest 会启动完整容器,不写就注入不了 RabbitTemplate。
  • @RabbitListener 由 RabbitListenerAnnotationBeanPostProcessor 包装成 MessageListener,容器启动后建连接、订阅队列,消息到了由容器自己的线程回调方法。参数类型决定反序列化成什么,发什么类型就得用什么类型接。
  • 队列必须先存在,不然消费者启动就报错,这个问题交给 3.4 的自动声明解决。
  • 方法抛异常默认会 requeue 重试,可能死循环刷日志。
  • 消费者线程不是 Web 请求线程,UserContext 里的 ThreadLocal 取不到值。

3.2 WorkQueue:能者多劳

核心是把多个消费者绑到同一个队列,一起消费。适用于消息处理耗时、生产速度大于消费速度的场景。

// 发送:连续发 50 条模拟堆积
for (int i = 0; i < 50; i++) {
    rabbitTemplate.convertAndSend("work.queue", "hello, message_" + i);
    Thread.sleep(20);
}
// 接收:两个消费者监听同一个队列
@RabbitListener(queues = "work.queue")
public void listenWorkQueue1(String msg) throws InterruptedException {
    System.out.println("消费者1接收到消息:【" + msg + "】");
    Thread.sleep(20);   // 模拟每秒处理 50 条
}
 
@RabbitListener(queues = "work.queue")
public void listenWorkQueue2(String msg) throws InterruptedException {
    System.err.println("消费者2接收到消息:【" + msg + "】");
    Thread.sleep(200);  // 模拟每秒处理 5 条
}

按默认设置跑,两个消费者各拿 25 条。队列是轮询分发,一条一条轮流给,不看消费者处理能力,快的干完就空闲,慢的还在慢慢磨。

在 consumer 的 yml 里加一行:

spring:
  rabbitmq:
    listener:
      simple:
        prefetch: 1 # 处理完一条并 ack 之后,才拿下一批

再跑一次,消费者1 处理 44 条、消费者2 只有 6 条,总耗时还是一秒左右。

prefetch 是消费者的预取数量,指未 ack 之前最多同时持有几条未确认消息。Spring AMQP 的 simple 容器默认 250,等于一次把一堆消息揽到怀里,队列以为它很能干就不再给别人;设成 1 就是处理完再拿,队列能把新消息派给空闲的消费者。

小结:多个消费者绑同一个队列时,同一条消息只会被一个消费者处理;用 prefetch 控制预取数量来实现能者多劳。

3.3 交换机与三种路由模型

引入交换机后,publisher 不再发给队列而是发给交换机,queue 必须和交换机绑定,consumer 不变。交换机只负责转发,没有队列绑定或者没有符合规则的队列,消息就丢了。

类型路由规则
Fanout广播,投给所有绑定到它的队列
Direct按 RoutingKey 完全一致匹配
TopicRoutingKey 支持通配符,比 Direct 灵活
Headers按消息头匹配,用得少

Fanout

// 发送:Fanout 忽略 routingKey,传 "" 即可
// 不能传 null,会抛 IllegalArgumentException
rabbitTemplate.convertAndSend("hmall.fanout", "", "hello, everyone!");
@RabbitListener(queues = "fanout.queue1")
public void listenFanoutQueue1(String msg) {
    System.out.println("消费者1接收到Fanout消息:【" + msg + "】");
}
 
@RabbitListener(queues = "fanout.queue2")
public void listenFanoutQueue2(String msg) {
    System.out.println("消费者2接收到Fanout消息:【" + msg + "】");
}

一条消息两个消费者都收到,这点和 WorkQueue 正相反。

Direct

队列绑定时要指定 bindingKey,发送时要指定 routingKey,两者完全一致队列才收。假设 direct.queue1 绑了 blue 和 red,direct.queue2 绑了 yellow 和 red:

rabbitTemplate.convertAndSend("hmall.direct", "red", "红色警报!");
// red 两个队列都绑了 → 两个消费者都收到
 
rabbitTemplate.convertAndSend("hmall.direct", "blue", "虚惊一场!");
// blue 只有 direct.queue1 绑了 → 只有消费者1收到

多个队列用同一个 key 时,效果就和 Fanout 差不多了。

Topic

同样按 RoutingKey 路由,区别是绑定用的 BindingKey 可以写通配符。BindingKey 由多个单词组成,用 . 分隔。

通配符含义
#匹配零个或多个词
*匹配恰好一个词

item.# 能匹配 item.spu.insert 和 item.spu,item.* 只能匹配 item.spu。

假设 routingKey 有 china.news、china.weather、japan.news、japan.weather:topic.queue1 绑 china.#,能收到两条中国的;topic.queue2 绑 #.news,能收到两条新闻。发一条 china.news,两个队列都匹配,两个消费者都收到。

rabbitTemplate.convertAndSend("hmall.topic", "china.news", "喜报!孙悟空大战哥斯拉,胜!");

用哪个交换机看需求:消息要无条件被多个队列收到就 Fanout;有条件、但条件是几个固定 key 就 Direct;条件带层级、需要用通配符就 Topic。

3.4 声明队列和交换机

控制台手动建的问题是,队列和交换机由程序员定义,上线却要运维创建,中间靠文档传递,很容易出错。推荐让程序启动时检查并自动创建。

三种方式:控制台手动(不推荐);配置类里声明 Queue、Exchange、Binding 三个 Bean(直观但啰嗦);@RabbitListener 配合 @QueueBinding 注解声明(最简洁,推荐)。

配置类方式:

@Configuration
public class FanoutConfig {
    @Bean
    public FanoutExchange fanoutExchange() {
        return new FanoutExchange("hmall.fanout");
    }
 
    @Bean
    public Queue fanoutQueue1() {
        return new Queue("fanout.queue1");
    }
 
    @Bean
    public Binding bindingQueue1(Queue fanoutQueue1, FanoutExchange fanoutExchange) {
        return BindingBuilder.bind(fanoutQueue1).to(fanoutExchange);
    }
    // 第二个队列同理
}

RabbitAdmin 会扫描容器里所有 Queue、Exchange、Binding 类型的 Bean,连接建立后自动向 Broker 声明,所以 Bean 名字不重要,类型才重要。Binding 方法的参数是按类型注入的,参数名只影响可读性。Fanout 绑定不需要 key,bind().to() 就结束了。

Direct 的痛点在于每个 key 都要写一个 Binding Bean,4 个 key 就是 4 个 Bean,注解方式就是为解决这个出现的:

@RabbitListener(bindings = @QueueBinding(
    value = @Queue(name = "direct.queue1"),
    exchange = @Exchange(name = "hmall.direct", type = ExchangeTypes.DIRECT),
    key = {"red", "blue"}
))
public void listenDirectQueue1(String msg) {
    System.out.println("消费者1接收到direct.queue1的消息:【" + msg + "】");
}
 
@RabbitListener(bindings = @QueueBinding(
    value = @Queue(name = "topic.queue1"),
    exchange = @Exchange(name = "hmall.topic", type = ExchangeTypes.TOPIC),
    key = "china.#"
))
public void listenTopicQueue1(String msg) {
    System.out.println("消费者1接收到topic.queue1的消息:【" + msg + "】");
}

两个方法就干完了原来一整个配置类的活。

常用属性速查:

注解属性说明
@Queuename / value队列名,两者等价
durable是否持久化,默认 true
exclusive是否排他,默认 false
autoDelete没有消费者时是否自动删除,默认 false
arguments额外参数,如死信、TTL、队列长度限制
@Exchangename / value交换机名
typeExchangeTypes.DIRECT / TOPIC / FANOUT / HEADERS
durable是否持久化,默认 true
delayed是否延迟交换机,需要插件
@QueueBindingkey绑定 key,数组形式可以写多个

如果 Broker 上已经存在同名但属性不同的队列(比如原来 durable=false,现在声明成 true),启动会报 406 PRECONDITION_FAILED,先删掉旧队列或者把属性对齐。

3.5 消息转换器

RabbitTemplate 收的是 Object,发送时序列化成字节,接收时再反序列化回来,默认用 JDK 序列化。问题是数据体积大、有反序列化攻击风险、控制台里看到的是一堆 rO0ABXNy... 乱码。

改用 JSON。两端都要加依赖,如果项目里引了 spring-boot-starter-web 就不用再加,它已经带了 Jackson:

<dependency>
    <groupId>com.fasterxml.jackson.dataformat</groupId>
    <artifactId>jackson-dataformat-xml</artifactId>
    <version>2.9.10</version>
</dependency>

再在两端的启动类里加 Bean:

@Bean
public MessageConverter messageConverter() {
    Jackson2JsonMessageConverter converter = new Jackson2JsonMessageConverter();
    // 自动生成消息 id,将来做幂等判断时能识别重复消息
    converter.setCreateMessageIds(true);
    return converter;
}

Bean 的类型必须写成 MessageConverter,自动装配会把它同时用到 RabbitTemplate 和监听容器工厂上,收发都是 JSON。两端都要配,一端 JDK 一端 JSON 会直接反序列化失败。发送方发 Map,接收方就得用 Map<String, Object> 接;发自定义对象则要有无参构造器。

@RabbitListener(queues = "object.queue")
public void listenObjectQueue(Map<String, Object> msg) {
    System.out.println("消费者接收到object.queue消息:【" + msg + "】");
}

Jackson2JsonMessageConverter 真正依赖的是 jackson-databind,上面的 jackson-dataformat-xml 只是把它传递依赖带进来了,直接引 jackson-databind 更规范。


四、业务改造:支付成功异步通知交易服务

把”支付成功后用 OpenFeign 调交易服务改订单状态”改成”发一条 MQ 消息”。

项值
交换机pay.topic(topic 类型)
队列mark.order.pay.queue
BindingKey / RoutingKeypay.success
消息内容订单 id

消费端(trade-service):

@Component
@RequiredArgsConstructor
public class PayStatusListener {
 
    private final IOrderService orderService;
 
    @RabbitListener(bindings = @QueueBinding(
            value = @Queue(name = "mark.order.pay.queue", durable = "true"),
            exchange = @Exchange(name = "pay.topic", type = ExchangeTypes.TOPIC),
            key = "pay.success"
    ))
    public void listenPaySuccess(Long orderId) {
        orderService.markOrderPaySuccess(orderId);
    }
}

队列、交换机、绑定三种声明一行注解全搞定。durable = "true" 表示队列持久化,Broker 重启后还在,持久化队列配持久化消息才真不丢。生产端发 Long,接收参数就必须是 Long。

生产端(pay-service)在 tryPayOrderByBalance 里把原来的 Feign 调用换掉:

try {
    rabbitTemplate.convertAndSend("pay.topic", "pay.success", po.getBizOrderNo());
} catch (Exception e) {
    log.error("支付成功的消息发送失败,支付单id:{},交易单id:{}", po.getId(), po.getBizOrderNo(), e);
}

改完的差别:从强依赖交易服务可用变成只依赖 MQ 可用;发消息近乎瞬时,不用等下游响应;交易服务挂了消息还在队列里等着,恢复后接着消费,不会再回滚支付。

这段代码留了两个基础篇式的妥协,高级篇会解决:

一是 try/catch 把发送异常吞掉了。为了不让通知失败回滚支付这个核心业务,这么做是对的,代价是消息真发不出去时没人知道,订单会一直停在待支付。高级篇用发送方确认加重试,或者延迟消息兜底。

二是发消息在事务提交之前。消息发出去之后事务回滚,就会出现交易服务改了订单状态、支付单没改的不一致。规范做法是用事务同步器在提交后再发消息,或者用本地消息表。

实际 hmall 代码里 PayStatusListener 还做了幂等处理,把无条件更新换成了带 status = 1 条件的更新,重复消费也不会出错。


五、练习要点

MQ 配置抽取到 Nacos:每个服务的 bootstrap.yaml 里在 shared-configs 加一条 shared-mq.yaml,内容就是那段 spring.rabbitmq.*。

下单后异步清理购物车:交换机 trade.topic,队列 cart.clear.queue,key 是 order.create,消息体是商品 id 集合。

try {
    rabbitTemplate.convertAndSend(
            MqConstants.TRADE_EXCHANGE_NAME, MqConstants.ORDER_CREATE_KEY, itemIds);
} catch (AmqpException e) {
    log.error("清理购物车的消息发送异常", e);
}
@RabbitListener(bindings = @QueueBinding(
        value = @Queue(name = "cart.clear.queue"),
        exchange = @Exchange(name = MqConstants.TRADE_EXCHANGE_NAME, type = ExchangeTypes.TOPIC),
        key = MqConstants.ORDER_CREATE_KEY
))
public void listenOrderCreate(List<Long> itemIds) {
    cartService.removeByItemIds(itemIds);
}

消息体是 List<Long>,接收参数也得是 List<Long>。

登录用户信息怎么传:异步调用不走 Web 请求线程,UserContext 里的 ThreadLocal 取不到值。笨办法是在消息体里带上用户、消费者再取出来,麻烦而且体验不统一。优雅做法是用 MessagePostProcessor 往消息 header 里塞:

// 发送端:发出前把用户写进 header
rabbitTemplate.setBeforePublishPostProcessors(message -> {
    Long userId = UserContext.getUser();
    if (userId != null) {
        message.getMessageProperties().setHeader("user-info", userId);
    }
    return message;
});
 
// 接收端:收到后从 header 取出来放回 UserContext
rabbitListenerContainerFactory.setAfterReceivePostProcessors(message -> {
    Long userId = message.getMessageProperties().getHeader("user-info");
    if (userId != null) {
        UserContext.setUser(userId);
    }
    return message;
});

配好之后业务代码里照旧 UserContext.getUser(),用 MQ 的人完全无感。

哪些业务适合改异步:非核心链路(短信、积分、日志)、实时性要求低、能接受最终一致性、一个事件触发多个下游动作。需要立刻拿返回值做判断的(比如扣库存前要先知道够不够)和强一致的不能改。


六、速查

概念类:

要点说明
交换机不存储消息没有匹配的队列,消息直接丢
队列必须绑定交换机只建队列不绑定,消息到不了
默认交换机名字是 "",每个队列自动绑定,bindingKey 等于队列名
virtual host数据隔离单位,yml 里配 virtual-host
端口5672 通信,15672 控制台

代码易错点:

易错点正确做法
yml port 写成 15672通信端口是 5672
convertAndSend 传 null routingKey传 ""
发 Long 收 String类型必须一致
只在一端配 JSON 转换器两端都要配
自定义对象接收报错对象要有无参构造器
@Bean 声明队列没生效检查 Bean 类型是不是 Queue / Exchange / Binding
注解声明报 406Broker 上已有同名但属性不同的队列,先删掉
Direct 写了一堆 Binding Bean改用 @QueueBinding

WorkQueue:

要点说明
同一条消息只被一个消费者处理这是和 Fanout 的本质区别
默认轮询分发不看处理能力,快的空闲慢的堆积
prefetch: 1处理完一条再拿一条,能者多劳
默认 prefetchSpring AMQP simple 容器是 250

基础篇留下的问题,正是高级篇的内容:

基础篇写法遗留问题高级篇解法
try/catch 吞掉发送异常消息可能真丢了,无人知晓发送方确认 + 重试
消费者抛异常默认 requeue可能死循环刷日志消费者确认 + 失败重试 + 死信队列
事务中直接发消息事务回滚与消息不一致事务同步器 / 本地消息表
队列和消息持久化极端情况仍可能丢惰性队列、镜像队列、仲裁队列
订单超时自动取消基础篇没讲延迟消息 + 延迟交换机插件
消息重复消费没有幂等保障幂等设计(messageId、状态条件更新)
登录用户信息传递传递麻烦MessagePostProcessor 传 header

精简自课程文档《MQ基础》(原文 1307 行),保留核心概念、五种模型、关键代码和易错点。