1、RabbitMQ环境搭建和上手体验
本文介绍了windows环境下的单机
RabbitMQ的搭建和docker、k8s环境下的RabbitMQ集群搭建,最后通过一个SpringBoot+RabbitMQ集成项目简单演示了RabbitMQ的基本使用
1、RabbitMQ环境的搭建
由于rabbitmq是由erlang语言编写的,所以,如果要在服务器(windows或Linux等)上直接安装rabbitmq,则需要先在对应服务器上安装erlang环境。而且,erlang环境的版本和rabbitmq的版本号还有一个对应关系,如果版本不匹配,则rabbitmq可能无法正常运行,会报一个不匹配的错。具体的版本对应关系参考:
我这里选择的是4.2.0-rc.1,所以对应的erlang版本选择的是OTP-27.3.4.14
erlang:https://github.com/erlang/otp/releases#release-OTP-27.3.4.13
rabbitmq:https://github.com/rabbitmq/rabbitmq-server/releases?page=2#release-v4.2.0-rc.1
1.1、Windows单机环境
1.1.1、erlang安装
- 安装包解压:

- 配置环境变量
- 验证

如果没有报错,能够正常显示出erlang的版本号,则说明erlang环境配置成功
1.1.2、安装rabbitmq
- 安装包解压

其中sbin目录就是操作rabbitmq的所有命令工具:

- 启动
解压之后什么都不用做,通过命令就可以直接运行
# 这里先启用rabbitmq的管理插件
rabbitmq-plugins.bat enable rabbitmq_management
# 正式启动
rabbitmq-server.bat start
运行结果:

看到这个completed就说明启动成功了,我们可以打开默认的管理界面看一下:
默认的账号和密码都是
guest

- 配置环境变量
虽然我们成功启动了rabbitmq,但是每次操作还是要先进入到安装目录的sbin目录下才能执行命令,为了方便后续操作,可以将sbin目录配置到系统PATH中
- 服务化
rabbitmq提供了一个命令,可以将rabbitmq注册成系统的一个服务,这样就可以让操作系统去调度,比如服务器重启后自动运行,这个命令的说明在解压后根目录的readme-service.txt,它告诉我们通过rabbitmq-service.bat install就可以将rabbitmq注册到系统服务上去,它告诉我们通过rabbitmq-service.bat enable就可以让RabbitMQ跟随系统服务进行管理
# 注意,我这里已经配置好了rabbitmq的环境变量,否则就要指定到rabbitmq-service.bat而不是直接用rabbitmq-service
D:\tools\rabbitmq\rabbitmq_server-4.2.0-rc.1>rabbitmq-service install
D:\tools\erlang\otp_win64_27.3.4.14\erts-15.2.7.10\bin\erlsrv: Service RabbitMQ added to system.
服务注册成功之后,就可以通过这里提示的服务名称来启动或停止服务了
net start RabbitMQ # 启动rabbitmq
net stop RabbitMQ # 停止rabbitmq

1.1.3、安装插件(以安装延迟消息插件为例)
- 下载插件(如果rabbitmq本身没有自带),放到
plugins目录下

- 启用插件
rabbitmq-plugins enable rabbitmq_delayed_message_exchange

1.1.4、数据存储
我们知道rabbitmq是典型的MQ,采用的是队列来存储消息数据,但是最终的数据持久化是如何实现的呢?其实rabbitmq采用的是一个文件类型数据库mnesia来存储的,默认的文件目录在:C:\Users\{用户名}\AppData\Roaming\RabbitMQ

在进到db目录就是mnesia的文件存储规则了

Linux环境下的安装步骤相同,只是安装包可能不同,这里不再赘述
1.2、docker集群环境
1.2.1、创建专属网络
docker network create rabbitmqtest
1.2.2、运行RabbitMQ的docker容器
docker run -d \
--name=rabbitmq1 \
-p 5673:5672 \
-p 15673:15672 \
-e RABBITMQ_NODENAME=rabbitmq1 \
-e RABBITMQ_ERLANG_COOKIE='hellorabbitmq' \
-h rabbitmq1 \
--net=rabbitmqtest \
registry.cn-hangzhou.aliyuncs.com/xhsx/rabbitmq:4.2.0-rc.1-management
docker run -d \
--name=rabbitmq2 \
-p 5674:5672 \
-p 15674:15672 \
-e RABBITMQ_NODENAME=rabbitmq2 \
-e RABBITMQ_ERLANG_COOKIE='hellorabbitmq' \
-h rabbitmq2 \
--net=rabbitmqtest \
registry.cn-hangzhou.aliyuncs.com/xhsx/rabbitmq:4.2.0-rc.1-management
docker run -d \
--name=rabbitmq3 \
-p 5675:5672 \
-p 15675:15672 \
-e RABBITMQ_NODENAME=rabbitmq3 \
-e RABBITMQ_ERLANG_COOKIE='hellorabbitmq' \
-h rabbitmq3 \
--net=rabbitmqtest \
registry.cn-hangzhou.aliyuncs.com/xhsx/rabbitmq:4.2.0-rc.1-management
运行成功之后,用默认的guest/guest账号登录管理控制台

1.2.3、创建集群
进入第二、三个容器,并将其以内存节点的角色加入到第一个容器里边,工程组成3个节点的集群
docker exec -it rabbitmq2 /bin/bash
rabbitmqctl stop_app
rabbitmqctl reset
rabbitmqctl join_cluster rabbitmq1@rabbitmq1
rabbitmqctl start_app
docker exec -it rabbitmq3 /bin/bash
rabbitmqctl stop_app
rabbitmqctl reset
rabbitmqctl join_cluster rabbitmq1@rabbitmq1
rabbitmqctl start_app
执行完成后,再看集群节点状态,可以看大有三个节点了

1.3、k8s集群环境
1.2.1、资源文件
注:这里所用到的命名空间是middleware,可以根据自己的情况酌情修改
1.2.1.1、configMap配置
cm.yaml
apiVersion: v1
kind: ConfigMap
metadata:
name: rmq-cluster-config
namespace: middleware
labels:
addonmanager.kubernetes.io/mode: Reconcile
data:
enabled_plugins: |
[rabbitmq_management,rabbitmq_peer_discovery_k8s].
rabbitmq.conf: |
loopback_users.guest = false
## Clustering
cluster_formation.peer_discovery_backend = rabbit_peer_discovery_k8s
cluster_formation.k8s.host = kubernetes.default.svc.cluster.local
cluster_formation.k8s.address_type = hostname
#################################################
# middleware is rabbitmq-cluster's namespace#
#################################################
cluster_formation.k8s.hostname_suffix = .rmq-cluster.middleware.svc.cluster.local
cluster_formation.node_cleanup.interval = 10
cluster_formation.node_cleanup.only_log_warning = true
cluster_partition_handling = autoheal
## queue master locator
queue_master_locator=min-masters
enabled_plugins:默认启用的插件
1.2.1.2、secret
secret.yaml:
apiVersion: v1
kind: Secret
metadata:
name: rmq-cluster-secret
namespace: middleware
stringData:
cookie: ERLANG_COOKIE
username: admin
password: abc12345!@@@
type: Opaque
1.2.1.3、service服务
svc.yaml
apiVersion: v1
kind: Service
metadata:
name: rmq-cluster
namespace: middleware
labels:
app: rmq-cluster
spec:
selector:
app: rmq-cluster
clusterIP: 10.96.0.11 #指定clusterIP,也可以用None不指定,访问的时候直接用serviceName+端口号
ports:
- name: http
port: 15672
protocol: TCP
targetPort: 15672
- name: amqp
port: 5672
protocol: TCP
targetPort: 5672
type: ClusterIP
1.2.1.4、pvc存储
注:这里用了nfs的storageClass动态分配存储空间,sc的定义如下:
sc.yaml:
apiVersion: storage.k8s.io/v1
kind: StorageClass
metadata:
namespace: middleware
name: sx-managed-nfs-storage
provisioner: sx-nfs-client-provisioner # or choose another name, must match deployment's env PROVISIONER_NAME'
parameters:
archiveOnDelete: "false"
---
kind: Deployment
apiVersion: apps/v1
metadata:
namespace: middleware
name: sx-nfs-client-provisioner
spec:
replicas: 1
strategy:
type: Recreate
selector:
matchLabels:
app: nfs-client-provisioner
template:
metadata:
namespace: middleware
labels:
app: nfs-client-provisioner
spec:
serviceAccount: nfs-client-provisioner
containers:
- name: sx-nfs-client-provisioner
imagePullPolicy: IfNotPresent
image: registry.cn-hangzhou.aliyuncs.com/xhsx/nfs-subdir-external-provisioner:v4.0.0
volumeMounts:
- name: sx-nfs-client-root
mountPath: /persistentvolumes
env:
- name: PROVISIONER_NAME
value: sx-nfs-client-provisioner
- name: NFS_SERVER
value: 192.168.0.1
- name: NFS_PATH
value: /nfs_win/k8s
volumes:
- name: sx-nfs-client-root
nfs:
server: 192.168.0.1
path: /nfs_win/k8s
这里用到了一个serviceAccount,定义如下:
serviceAccount.yaml:
---
kind: ClusterRole
apiVersion: rbac.authorization.k8s.io/v1
metadata:
namespace: middleware
name: nfs-client-provisioner-runner
rules:
- apiGroups: [""]
resources: ["persistentvolumes"]
verbs: ["get", "list", "watch", "create", "delete"]
- apiGroups: [""]
resources: ["persistentvolumeclaims"]
verbs: ["get", "list", "watch", "update"]
- apiGroups: ["storage.k8s.io"]
resources: ["storageclasses"]
verbs: ["get", "list", "watch"]
- apiGroups: [""]
resources: ["events"]
verbs: ["create", "update", "patch"]
---
kind: ClusterRoleBinding
apiVersion: rbac.authorization.k8s.io/v1
metadata:
namespace: middleware
name: run-nfs-client-provisioner
subjects:
- kind: ServiceAccount
name: nfs-client-provisioner
namespace: middleware
roleRef:
kind: ClusterRole
name: nfs-client-provisioner-runner
apiGroup: rbac.authorization.k8s.io
---
kind: Role
apiVersion: rbac.authorization.k8s.io/v1
metadata:
namespace: middleware
name: leader-locking-nfs-client-provisioner
rules:
- apiGroups: [""]
resources: ["endpoints"]
verbs: ["get", "list", "watch", "create", "update", "patch"]
---
kind: RoleBinding
apiVersion: rbac.authorization.k8s.io/v1
metadata:
namespace: middleware
name: leader-locking-nfs-client-provisioner
subjects:
- kind: ServiceAccount
name: nfs-client-provisioner
# replace with namespace where provisioner is deployed
namespace: middleware
roleRef:
kind: Role
name: leader-locking-nfs-client-provisioner
apiGroup: rbac.authorization.k8s.io
---
kind: ServiceAccount
apiVersion: v1
metadata:
namespace: middleware
name: nfs-client-provisioner
pvc.yaml:
---
apiVersion: v1
kind: PersistentVolumeClaim
metadata:
name: rabbitmq-cluster-storage
namespace: middleware
spec:
accessModes:
- ReadWriteMany
resources:
requests:
storage: 10Gi
storageClassName: sx-managed-nfs-storage
1.2.1.5、statefulSet
sts.yaml:
apiVersion: apps/v1
kind: StatefulSet
metadata:
name: rmq-cluster
namespace: middleware
labels:
app: rmq-cluster
spec:
replicas: 3
selector:
matchLabels:
app: rmq-cluster
serviceName: rmq-cluster
template:
metadata:
labels:
app: rmq-cluster
spec:
serviceAccountName: nfs-client-provisioner
terminationGracePeriodSeconds: 30
containers:
- name: rabbitmq
image: registry.cn-hangzhou.aliyuncs.com/xhsx/rabbitmq:3.7-management
imagePullPolicy: IfNotPresent
ports:
- containerPort: 15672
name: http
protocol: TCP
- containerPort: 5672
name: amqp
protocol: TCP
command:
- sh
args:
- -c
- cp -v /etc/rabbitmq/rabbitmq.conf ${RABBITMQ_CONFIG_FILE}; exec docker-entrypoint.sh
rabbitmq-server
env:
- name: TZ
value: Asia/Shanghai
- name: RABBITMQ_DEFAULT_USER
valueFrom:
secretKeyRef:
key: username
name: rmq-cluster-secret
- name: RABBITMQ_DEFAULT_PASS
valueFrom:
secretKeyRef:
key: password
name: rmq-cluster-secret
- name: RABBITMQ_ERLANG_COOKIE
valueFrom:
secretKeyRef:
key: cookie
name: rmq-cluster-secret
- name: K8S_SERVICE_NAME
value: rmq-cluster
- name: POD_IP
valueFrom:
fieldRef:
fieldPath: status.podIP
- name: POD_NAME
valueFrom:
fieldRef:
fieldPath: metadata.name
- name: POD_NAMESPACE
valueFrom:
fieldRef:
fieldPath: metadata.namespace
- name: RABBITMQ_USE_LONGNAME
value: "true"
- name: RABBITMQ_NODENAME
value: rabbit@$(POD_NAME).rmq-cluster.$(POD_NAMESPACE).svc.cluster.local
- name: RABBITMQ_CONFIG_FILE
value: /var/lib/rabbitmq/rabbitmq.conf
livenessProbe:
exec:
command:
- rabbitmqctl
- status
initialDelaySeconds: 30
timeoutSeconds: 10
readinessProbe:
exec:
command:
- rabbitmqctl
- status
initialDelaySeconds: 10
timespec.template.spec.outSeconds: 10
volumeMounts:
- name: config-volume
mountPath: /etc/rabbitmq
readOnly: false
- name: rabbitmq-storage
mountPath: /var/lib/rabbitmq
readOnly: false
volumes:
- name: config-volume
configMap:
items:
- key: rabbitmq.conf
path: rabbitmq.conf
- key: enabled_plugins
path: enabled_plugins
name: rmq-cluster-config
- name: rabbitmq-storage
persistentVolumeClaim:
claimName: rabbitmq-cluster-storage
注意几个配置要和上面定义的相同:
- spec.template.spec.volumes.persistentVolumeClaim.claimName
- spec.template.spec.volumes.configMap.name
- spec.template.spec.serviceAccountName
- spec.serviceName
按照依赖顺序分别apply一下即可完成rabbitmq集群的构建
kubectl apply -f cm.yaml
kubectl apply -f secret.yaml
kubectl apply -f svc.yaml
kubectl apply -f serviceAccount.yaml
kubectl apply -f sc.yaml
kubectl apply -f pvc.yaml
kubectl apply -f sts.yaml
1.2.1.6、访问
- 管理后台界面
虽然经过上面的步骤搭建了一个rabbitmq集群,但是如果是在k8s集群之外访问,则需要暴露一个nodePort类型的service出来,我这里是用了ingress-nginx网络
ingress.yaml:
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
namespace: middleware
name: rabbitmq-ingress
spec:
ingressClassName: nginx
rules:
- host: rabbitmq.shengxiao.com
http:
paths:
- path: /
pathType: Prefix
backend:
service:
name: rmq-cluster
port:
number: 15672
ubuntu@k8s-master-01:~/k8s/rabbitmq$ sudo kubectl -n ingress-nginx get po -o wide
NAME READY STATUS RESTARTS AGE IP NODE NOMINATED NODE READINESS GATES
ingress-nginx-admission-create-s4b4c 0/1 Completed 0 19d <none> k8s-worker-03 <none> <none>
ingress-nginx-admission-patch-grl6l 0/1 Completed 0 19d <none> k8s-worker-01 <none> <none>
ingress-nginx-controller-74b49756b9-q4npk 1/1 Running 2 (15d ago) 19d 192.168.0.141 k8s-worker-01 <none> <none>
可以看到我的ingress-nginx-controller运行在了192.168.0.141这个节点上的,所以我配置了自定义域名rabbitmq.shengxiao.com的本地映射,具体就是在我要访问rabbitmq的机器的hosts文件中添加一行记录:
192.168.0.141 rabbitmq.shengxiao.com
最后访问这个自定义域名访问集群(默认的账号密码:guest / guest):

- 应用访问
上面演示了如何访问rabbitmq的后台界面,但是在我们的应用代码中,如果想要往broker发送消息或者监听消息,通过的端口号是5672,这个时候如果是在k8s集群外部,就需要暴露一个nodePort,具体可以参考下面定义:
nodePort.yaml:
apiVersion: v1
kind: Service
metadata:
labels:
app: rmq-cluster
name: rabbitmq-svc-nodeport
namespace: middleware
spec:
type: NodePort
ports:
- port: 5672
protocol: TCP
targetPort: 5672
nodePort: 30672
selector:
app: rmq-cluster
这样,在配置rabbitmq的配置的时候就可以这样配置:
spring.rabbitmq.host=192.168.0.141
spring.rabbitmq.port=30672
spring.rabbitmq.virtual-host=/
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
spring.rabbitmq.host可以随便配置一个k8s的节点IP
当然,如果是在k8s集群内部的应用访问,则可以直接通过serviceName来访问,这时候rabbitmq的配置文件就可以改成这样:
spring.rabbitmq.addresses=rmq-cluster.middleware:5672
spring.rabbitmq.virtual-host=/
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
rmq-cluster:前面定义的serviceName
middleware:命名空间名称
5672:定义service的时候的service端口号
2、RabbitMQ上手体验
SpringBoot集成
2.1、添加依赖
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<spring.boot.version>2.7.18</spring.boot.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-logging</artifactId>
</dependency>
</dependencies>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-dependencies</artifactId>
<version>${spring.boot.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
2.2、配置文件
2.2.1、springboot配置文件
spring.rabbitmq.addresses=localhost:5672
spring.rabbitmq.virtual-host=test
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
# 日志配置文件路径
logging.config=classpath:logback.xml
2.2.2、日志配置文件
<configuration scan="true" scanPeriod="60 seconds" debug="false">
<appender name="console" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<!--格式化输出:%d表示日期,%thread表示线程名,%-5level:级别从左显示5个字符宽度,%msg:日志消息,%n是换行符-->
<pattern>%d{yyyy-MM-dd HH:mm:ss,SSS} :[ %-5level] %thread %logger{50} - %msg%n</pattern>
</encoder>
</appender>
<root level="INFO">
<appender-ref ref="console"/>
</root>
</configuration>
2.3、代码
2.3.1、配置类
package com.study.rabbitmq.config;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.amqp.core.*;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class RabbitMQConfiguration {
private static final Logger logger = LoggerFactory.getLogger(RabbitMQConfiguration.class);
// 直连类型交换机
@Bean("directExchange")
public DirectExchange getDirectExchange(){
return new DirectExchange("test_direct_exchange");
}
// 第一个队列:绑定到直连交换机上
@Bean("directQueue")
public Queue getDirectQueue(){
Map<String, Object> args = new HashMap<String, Object>();
args.put("x-message-ttl",6000);
Queue queue = new Queue("direct_queue", false, false, true, args);
return queue;
}
// 创建交换机-队列绑定关系
@Bean
public Binding bindDirectQueue(@Qualifier("directQueue") Queue queue, @Qualifier("directExchange") DirectExchange exchange){
return BindingBuilder.bind(queue).to(exchange).with("points.change");
}
// 主题类型交换机
@Bean("topicExchange")
public TopicExchange getTopicExchange(){
return new TopicExchange("test_topic_exchange");
}
@Bean("topicFirstQueue")
public Queue getTopicFirstQueue(){
return new Queue("topic_first_queue");
}
@Bean("topicSecondQueue")
public Queue getTopicSecondQueue(){
return new Queue("topic_second_queue");
}
// 创建交换机-队列绑定关系
@Bean
public Binding bindTopicFirstQueue(@Qualifier("topicFirstQueue") Queue queue, @Qualifier("topicExchange") TopicExchange exchange){
return BindingBuilder.bind(queue).to(exchange).with("order.#");
}
// 创建交换机-队列绑定关系
@Bean
public Binding bindTopicSecondQueue(@Qualifier("topicSecondQueue") Queue queue, @Qualifier("topicExchange") TopicExchange exchange){
return BindingBuilder.bind(queue).to(exchange).with("user.#");
}
// 广播类型交换机
@Bean("fanoutExchange")
public FanoutExchange getFanoutExchange(){
return new FanoutExchange("test_fanout_exchange");
}
@Bean("fanoutFirstQueue")
public Queue getfirstFanoutQueue(){
return new Queue("fanout_first_queue");
}
@Bean("fanoutSecondQueue")
public Queue getSecondFanoutQueue(){
return new Queue("fanout_second_queue");
}
@Bean
public Binding bindFanoutFirstQueue(@Qualifier("fanoutFirstQueue") Queue queue,@Qualifier("fanoutExchange") FanoutExchange exchange){
return BindingBuilder.bind(queue).to(exchange);
}
@Bean
public Binding bindFanoutSecondQueue(@Qualifier("fanoutSecondQueue") Queue queue,@Qualifier("fanoutExchange") FanoutExchange exchange){
return BindingBuilder.bind(queue).to(exchange);
}
}
注:
- 这里在代码里进行了交换机、队列以及它们之间的绑定关系,实际工作中应该由管理员在管理后台手动创建
- 这里先上手体验,后续再分析他们的工作原理、区别以和特性
- 这里我定义了3个交换机:一个直连交换机、一个主题类型交换机、一个广播类型交换机
- 直连类型交换机绑定了一个队列、主题类型绑定了两个队列、广播类型交换机绑定了两个队列
2.3.2、生产者
package com.study.rabbitmq.producer;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Component
public class RabbitMqProducer {
@Resource
private RabbitTemplate rabbitTemplate;
public void sendMessage(String exchange, String routingKey, String message){
this.rabbitTemplate.convertAndSend(exchange, routingKey, message);
}
}
2.3.3、消费者
上面定义了5个队列,这里定义5个消费者分别来消费这5个队列里的消息
DirectQueueConsumer.java:消费绑定直连交换机的队列
package com.study.rabbitmq.consumer;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
@RabbitListener(queues = "direct_queue")
public class DirectQueueConsumer {
@RabbitHandler
public void process(String msg){
System.out.println(" direct queue received msg : " + msg);
}
}
TopicFirstQueueConsumer.java:消费绑定第一个主题类型交换机的队列
package com.study.rabbitmq.consumer;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
@RabbitListener(queues = "topic_first_queue")
public class TopicFirstQueueConsumer {
@RabbitHandler
public void process(String msg){
System.out.println(" topic first queue received msg : " + msg);
}
}
TopicSecondQueueConsumer.java:消费绑定第二个主题类型交换机的队列
package com.study.rabbitmq.consumer;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
@RabbitListener(queues = "topic_second_queue")
public class TopicSecondQueueConsumer {
@RabbitHandler
public void process(String msg){
System.out.println(" topic second queue received msg : " + msg);
}
}
FanoutFirstQueueConsumer.java:消费绑定第一个广播类型交换机的队列
package com.study.rabbitmq.consumer;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
@RabbitListener(queues = "fanout_first_queue")
public class FanoutFirstQueueConsumer {
@RabbitHandler
public void process(String msg){
System.out.println(" fanout first queue received msg : " + msg);
}
}
FanoutSecondQueueConsumer.java:消费绑定第二个广播类型交换机的队列
package com.study.rabbitmq.consumer;
import org.springframework.amqp.rabbit.annotation.RabbitHandler;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
@RabbitListener(queues = "fanout_second_queue")
public class FanoutSecondQueueConsumer {
@RabbitHandler
public void process(String msg){
System.out.println(" fanout second queue received msg : " + msg);
}
}
2.3.4、定义一个接口,用来发送消息
package com.study.rabbitmq.controller;
import com.study.rabbitmq.bean.SendMessageRequest;
import com.study.rabbitmq.producer.RabbitMqProducer;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
@RestController
@RequestMapping("/message")
public class MessageController {
@Resource
private RabbitMqProducer rabbitMqProducer;
@PostMapping("/send")
public void sendMessage(@RequestBody SendMessageRequest request){
this.rabbitMqProducer.sendMessage(request.getExchange(),request.getRoutingKey(), request.getMessage());
}
}
2.3.5、启动类
package com.study.rabbitmq;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class RabbitMqBootstrap {
public static void main(String[] args) {
SpringApplication.run(RabbitMqBootstrap.class, args);
}
}
2.4、运行
2.4.1、启动
启动程序,可以在管理后台看到交换机和队列已经自动给我们创建好了,并且也自动建立了绑定关系
交换机:

队列:

队列绑定关系:

2.4.2、发送消息
2.4.2.1、直连交换机

由于我们在定义绑定关系的时候定义的routingKey是points.change,这里我们在发送消息的时候是随便写的,点击发送后发现消息已经成功被队列接收

但是消费端并没有接收到对应的消息,接下来我们改成正确的routingKey重新请求
{
"exchange": "test_direct_exchange",
"routingKey": "points.change",
"message": "我是一条发送到直连交换机里的消息"
}
可以看到消费端正确拿到了消息:

2.4.2.2、主题类型交换机
修改发送消息的请求体:
{
"exchange": "test_topic_exchange",
"routingKey": "order.created",
"message": "我是订单创建消息"
}
这里将交换机改成了主题类型的交换机,routingKey设置为“order.created”,看最终被哪个队列的消费者接收到,发送请求,查看代码控制台

可以看到消息被主题类型交换机绑定的第一个队列收到了,从上面的定义可以看到它们的绑定路由键是"order.#",而第二个队列的绑定路由键是"user.#",而我们在发送消息的时候的routingKey是"order.created",大概就知道他们的关系了吧?
注:这里还涉及一个
*和.的区别,只是不在这里阐述,后续再进行介绍
2.4.2.3、广播类型交换机
继续修改发送消息的请求体:
{
"exchange": "test_fanout_exchange",
"routingKey": "order.created",
"message": "我是发往广播类型交换机的消息"
}
这里只是改一下交换机的名称,看最终控制台打印出来的结果:

可以看到绑定到广播类型交换机的两个队列的消费者都收到了消息
