RabbitMQ是实现AMQP(高级消息队列协议)消息中间件的一种, kafuka是另外一种
RabbitMQ主要是为了实现系统之间的双向解耦而实现的,消息的发送者无需知道消息使用者的存在,反之亦然
当生产者大量产生数据时, 消费者无法快速消费, 那么需要一个中间层,保存这个数据
AMQP的主要特征是面向消息(Message)、队列(Queue)、路由(Exchange包括点对点和发布/订阅)、可靠性、安全
1.环境搭建
使用docker启动已经做好的镜像即可
docker pull rabbitmq:3-management
docker run -itd --rm --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management
2.依赖
pom.xml 增加依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
3.配置
在需要使用rabbitmq的项目配置文件中增加 rabbitmq 相关配置项
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
virtual-host: /
listener:
direct:
acknowledge-mode: manual
simple:
acknowledge-mode: manual
4.相关概念
一般来说消息队列有三个概念: 发消息者、队列、收消息者,
RabbitMQ 在这个基本概念之上, 多做了一层抽象, 在发消息者和队列之间, 加入了交换器 (Exchange).
这样发消息者和队列就没有直接联系, 转而变成发消息者把消息给交换器, 交换器根据调度策略再把消息再给队列。
交换机(Exchange)
交换机的功能主要是接收消息并且转发到绑定的队列,交换机不存储消息,在启用ack模式后,交换机找不到队列会返回错误。
交换机有四种类型:Direct, topic, Headers and Fanout
Direct:direct 类型的行为是“先匹配, 再投送”. 即在绑定时设定一个 routing_key, 消息的routing_key 匹配时, 才会被交换器投送到绑定的队列中去.
Topic:按规则转发消息(最灵活)
Headers:设置header attribute参数类型的交换机
Fanout:转发消息到所有绑定队列
项目中使用的是Topic模式
消费者模糊匹配路由Key,进行消息消费,路由Key必须是一串字符,用句号(.) 隔开,比如说,applet.create,applet.#
主要有两种模糊匹配:# 匹配一个或多个,* 匹配一个,一般使用#号匹配多个,*号用的比较少。
如果消费者和生产者在分属两个不同服务进程中,建议创建队列的操作放在消费者那个进程中做,因为消费者启动的时候需要监听队列,若放在了生产者
进程中创建而生产者又没有启动,则消费者启动的时候会报错。交换机在生产者和消费者都可以创建,就看谁先启动谁先创建了
5.代码
生产者进程相关代码
rabbitmq使用前的准备配置代码
@Configuration
public class RabbitmqConfig {
//使用Jackson2JsonMessageConverter 消息传递后转对象
@Bean
public MessageConverter messageConverter() {
return new Jackson2JsonMessageConverter();
}
//创建一个名为applet.exchange的交换机
@Bean
TopicExchange appletExchange() {
return new TopicExchange("applet.exchange", true, false);
}
}
发送消息到交换机代码
@Service
public class UpdaterServiceImpl implements UpdaterService {
@Autowired
RabbitTemplate rabbitTemplate;
public void create(AppletRequestDTO dto) {
rabbitTemplate.convertAndSend("applet.exchange", "applet.create", dto);
}
public void update(AppletRequestDTO dto) {
rabbitTemplate.convertAndSend("applet.exchange", "applet.update", dto);
}
}
消费者进程相关代码
rabbitmq使用前的准备配置代码
@Configuration
public class RabbitmqConfig {
//使用Jackson2JsonMessageConverter 消息传递后转对象
@Bean
public MessageConverter messageConverter() {
return new Jackson2JsonMessageConverter();
}
//创建一个名为applet.exchange的交换机
@Bean
TopicExchange appletExchange() {
return new TopicExchange("applet.exchange", true, false);
}
//创建一个名为search_index_queue队列
@Bean
Queue searchIndexQueue() {
return new Queue("search_index_queue", true);
}
//创建一个名为search_update_queue队列
@Bean
Queue searchUpdateQueue() {
return new Queue("search_update_queue", true);
}
//队列绑定到路由器上并指定使用什么路由Key
@Bean
Binding searchIndexBinding() {
return BindingBuilder.bind(searchIndexQueue()).to(appletExchange()).with("applet.create");
}
@Bean
Binding searchUpdateBinding() {
return BindingBuilder.bind(searchUpdateQueue()).to(appletExchange()).with("applet.update");
}
}
消费监听和接收消息队列的消息代码
有2种方式配置RabbitListener,一种使用RabbitListener标注类且使用RabbitHandler标注方法,另一种使用RabbitListener标注方法
@Component
@Slf4j
public class RabbitmqListener {
@Autowired
AppletSearchService appletSearchService;
@RabbitListener(queues = "search_index_queue")
public void handlerIndex(AppletRequestDTO dto, Channel channel, Message message) {
try {
//处理消息
appletSearchService.index(dto);
//通知消息处理完毕
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
} catch (Exception e) {
e.printStackTrace();
//丢弃这条消息
log.error("索引消息处理错误. msg:{}", e.getMessage());
}
}
@RabbitListener(queues = "search_update_queue")
public void handlerUpdate(AppletRequestDTO dto, Channel channel, Message message) {
try {
appletSearchService.update(dto);
//通知消息处理完毕
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
} catch (IOException e) {
e.printStackTrace();
//丢弃这条消息
//channel.basicNack(message.getMessageProperties().getDeliveryTag(), false,false);
log.error("更新文档消息消息处理错误. msg:{}", e.getMessage());
}
}
}