Skip to content

RabbitMQ

简介

RabbitMQ 是使用 Erlang 语言开发的开源消息代理软件,实现了高级消息队列协议(AMQP)。以稳定可靠、功能丰富、管理界面友好著称。

核心概念

  • Producer:消息生产者
  • Consumer:消息消费者
  • Exchange:交换机,消息路由的核心
  • Queue:队列,存储消息的缓冲区
  • Binding:绑定,定义 Exchange 和 Queue 之间的路由关系
  • Virtual Host:虚拟主机,多租户隔离

交换机类型

Direct Exchange(直连交换机)

根据 routing key 精确匹配:

java
channel.exchangeDeclare("direct-exchange", BuiltinExchangeType.DIRECT);
channel.queueBind("queue-1", "direct-exchange", "error");
channel.queueBind("queue-2", "direct-exchange", "info");

Topic Exchange(主题交换机)

通配符匹配 routing key:

java
channel.exchangeDeclare("topic-exchange", BuiltinExchangeType.TOPIC);
channel.queueBind("queue-1", "topic-exchange", "order.*");    // 匹配 order.create, order.pay
channel.queueBind("queue-2", "topic-exchange", "#.completed"); // 匹配任何以 completed 结尾

Fanout Exchange(广播交换机)

忽略 routing key,发送到所有绑定的队列:

java
channel.exchangeDeclare("fanout-exchange", BuiltinExchangeType.FANOUT);
channel.queueBind("queue-1", "fanout-exchange", "");
channel.queueBind("queue-2", "fanout-exchange", "");

消息确认机制

java
// 生产者确认
channel.confirmSelect();
if (channel.waitForConfirms()) {
    System.out.println("消息确认发送成功");
}

// 消费者手动 ACK
boolean autoAck = false;
channel.basicConsume("queue", autoAck, (consumerTag, delivery) -> {
    String message = new String(delivery.getBody(), "UTF-8");
    System.out.println("收到: " + message);
    channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
}, consumerTag -> {});

适用场景

  • 任务队列:异步处理耗时任务
  • 日志分发:将日志消息路由到不同处理系统
  • 微服务通信:服务间的异步消息传递
  • RPC 调用:基于 RabbitMQ 实现远程过程调用