您好,登錄后才能下訂單哦!
RabbitMQ簡(jiǎn)介
RabbitMQ是一個(gè)在AMQP基礎(chǔ)上完整的,可復(fù)用的企業(yè)消息系統(tǒng)
MQ全稱(chēng)為Message Queue, 消息隊(duì)列(MQ)是一種應(yīng)用程序?qū)?yīng)用程序的通信方法。應(yīng)用程序通過(guò)讀寫(xiě)出入隊(duì)列的消息(針對(duì)應(yīng)用程序的數(shù)據(jù))來(lái)通信,而無(wú)需專(zhuān)用連接來(lái)鏈接它們。消息傳遞指的是程序之間通過(guò)在消息中發(fā)送數(shù)據(jù)進(jìn)行通信,而不是通過(guò)直接調(diào)用彼此來(lái)通信,直接調(diào)用通常是用于諸如遠(yuǎn)程過(guò)程調(diào)用的技術(shù)。排隊(duì)指的是應(yīng)用程序通過(guò) 隊(duì)列來(lái)通信。隊(duì)列的使用除去了接收和發(fā)送應(yīng)用程序同時(shí)執(zhí)行的要求。
AMQP就是一個(gè)協(xié)議,是一個(gè)高級(jí)抽象層消息通信協(xié)議。
雖然在同步消息通訊的世界里有很多公開(kāi)標(biāo)準(zhǔn)(如 COBAR的 IIOP ,或者是 SOAP 等),但是在異步消息處理中卻不是這樣,只有大企業(yè)有一些商業(yè)實(shí)現(xiàn)(如微軟的 MSMQ ,IBM 的 Websphere MQ 等),因此,在 2006 年的 6 月,Cisco 、Redhat、iMatix 等聯(lián)合制定了 AMQP 的公開(kāi)標(biāo)準(zhǔn)。也就是說(shuō)AMQP是異步通訊的一個(gè)協(xié)議。
RabbitMQ使用場(chǎng)景
在項(xiàng)目中,將一些無(wú)需即時(shí)返回且耗時(shí)的操作提取出來(lái),進(jìn)行了異步處理,而這種異步處理的方式大大的節(jié)省了服務(wù)器的請(qǐng)求響應(yīng)時(shí)間,從而提高了系統(tǒng)的吞吐量。不過(guò)大多數(shù)不僅僅是無(wú)需即時(shí)返回,甚至是執(zhí)行是否成功都無(wú)所謂。如果需要即時(shí)返回則可以使用Dubbo,Spring boot與Dubbo集成可以去看Spring boot 集成Dubbox
RabbitMQ依賴
RabbitMQ并不是直接一個(gè)簡(jiǎn)單的jar包(Jar包只是提供一個(gè)基本的與RabbitMQ本身通訊的一些功能),和Dubbo相同,RabbitMQ也需要其他軟件來(lái)運(yùn)行,以下是RabbitMQ運(yùn)行所需要的軟件
1、Erlang
由于RabbitMQ軟件本身是基于Erlang開(kāi)發(fā)的,所以想要運(yùn)行RabbitMQ必須要先按照Erlang
Erlang官網(wǎng)
Erlang下載地址
RabbitMQ
RabbitMQ才是實(shí)現(xiàn)消息隊(duì)列的核心
RabbitMQ官網(wǎng)
RabbitMQ下載
配置RabbitMQ
安裝完成后,需要完成一些配置才能使用RabbitMQ,可以直接用cmd到RabbitMQ的安裝目錄下的sbin目錄通過(guò)命令配置,也可以直接在開(kāi)始菜單中直接找到RabbitMQ Command Prompt (sbin dir)運(yùn)行直接到達(dá)RabbitMQ的安裝目錄的sbin,為了方便,我們先啟用管理插件,執(zhí)行命令
rabbitmq-plugins.bat enable rabbitmq_management
即可,注意,這是在Windows下面,如果是Linux則沒(méi)有bat后綴 然后我們添加一個(gè)用戶,因?yàn)樵谕饩W(wǎng)環(huán)境沒(méi)有用戶的情況下是不能連接成功的,執(zhí)行添加用戶命令
rabbitmqctl.bat add_user springboot password
springboot是用戶名,password是密碼
然后為了方便演示,我們給springboot賦予管理員權(quán)限,方便登錄管理頁(yè)面
rabbitmqctl.bat set_user_tags springboot administrator
給賬號(hào)賦予虛擬主機(jī)權(quán)限
rabbitmqctl.bat set_permissions -p / springboot .* .* .*
然后啟動(dòng)RabbitMQ服務(wù) 訪問(wèn)RabbitMQ管理頁(yè)面http://localhost:15672即可看見(jiàn)登錄頁(yè)面,如果沒(méi)有創(chuàng)建用戶則可以用guest,guest登錄,如果有創(chuàng)建用戶則用創(chuàng)建的用戶登錄
創(chuàng)建Springboot項(xiàng)目
因?yàn)閯?chuàng)建spring boot項(xiàng)目在前面的文章已經(jīng)說(shuō)過(guò)很多次了,所以這里就不多說(shuō)了
添加RabbitMQ相關(guān)依賴
<!-- rabbitmq --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency>
沒(méi)錯(cuò),就是點(diǎn)配置,不過(guò)這樣可能有點(diǎn)不理解,我還是把全部配置貼出來(lái)吧
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>wang.raye.rabbitmq</groupId> <artifactId>demo1</artifactId> <version>0.0.1-SNAPSHOT</version> <packaging>jar</packaging> <name>demo1</name> <url>http://maven.apache.org</url> <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> </properties> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>1.4.0.RELEASE</version> </parent> <dependencies> <dependency> <groupId>junit</groupId> <artifactId>junit</artifactId> <version>3.8.1</version> <scope>test</scope> </dependency> <!-- Springboot --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <!-- rabbitmq --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> </dependencies> </project>
因?yàn)闆](méi)有做其他操作,所以目前項(xiàng)目主要是依賴2個(gè)模塊,一個(gè)Sprig boot,一個(gè)RabbitMQ
添加配置類(lèi)
package wang.raye.rabbitmq.demo1; import org.springframework.amqp.core.AcknowledgeMode; import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.DirectExchange; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.Queue; import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.core.ChannelAwareMessageListener; import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; /** * rabbitmq 的配置類(lèi) * * @author Raye * @since 2016年10月12日10:57:44 */ @Configuration public class RabbitMQConfig { /** 消息交換機(jī)的名字*/ public static final String EXCHANGE = "my-mq-exchange"; /** 隊(duì)列key1*/ public static final String ROUTINGKEY1 = "queue_one_key1"; /** 隊(duì)列key2*/ public static final String ROUTINGKEY2 = "queue_one_key2"; /** * 配置鏈接信息 * @return */ @Bean public ConnectionFactory connectionFactory() { CachingConnectionFactory connectionFactory = new CachingConnectionFactory("127.0.0.1",5672); connectionFactory.setUsername("springboot"); connectionFactory.setPassword("password"); connectionFactory.setVirtualHost("/"); connectionFactory.setPublisherConfirms(true); // 必須要設(shè)置 return connectionFactory; } /** * 配置消息交換機(jī) * 針對(duì)消費(fèi)者配置 FanoutExchange: 將消息分發(fā)到所有的綁定隊(duì)列,無(wú)routingkey的概念 HeadersExchange :通過(guò)添加屬性key-value匹配 DirectExchange:按照routingkey分發(fā)到指定隊(duì)列 TopicExchange:多關(guān)鍵字匹配 */ @Bean public DirectExchange defaultExchange() { return new DirectExchange(EXCHANGE, true, false); } /** * 配置消息隊(duì)列1 * 針對(duì)消費(fèi)者配置 * @return */ @Bean public Queue queue() { return new Queue("queue_one", true); //隊(duì)列持久 } /** * 將消息隊(duì)列1與交換機(jī)綁定 * 針對(duì)消費(fèi)者配置 * @return */ @Bean public Binding binding() { return BindingBuilder.bind(queue()).to(defaultExchange()).with(RabbitMQConfig.ROUTINGKEY1); } /** * 配置消息隊(duì)列2 * 針對(duì)消費(fèi)者配置 * @return */ @Bean public Queue queue1() { return new Queue("queue_one1", true); //隊(duì)列持久 } /** * 將消息隊(duì)列2與交換機(jī)綁定 * 針對(duì)消費(fèi)者配置 * @return */ @Bean public Binding binding1() { return BindingBuilder.bind(queue1()).to(defaultExchange()).with(RabbitMQConfig.ROUTINGKEY2); } /** * 接受消息的監(jiān)聽(tīng),這個(gè)監(jiān)聽(tīng)會(huì)接受消息隊(duì)列1的消息 * 針對(duì)消費(fèi)者配置 * @return */ @Bean public SimpleMessageListenerContainer messageContainer() { SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory()); container.setQueues(queue()); container.setExposeListenerChannel(true); container.setMaxConcurrentConsumers(1); container.setConcurrentConsumers(1); container.setAcknowledgeMode(AcknowledgeMode.MANUAL); //設(shè)置確認(rèn)模式手工確認(rèn) container.setMessageListener(new ChannelAwareMessageListener() { public void onMessage(Message message, com.rabbitmq.client.Channel channel) throws Exception { byte[] body = message.getBody(); System.out.println("收到消息 : " + new String(body)); channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); //確認(rèn)消息成功消費(fèi) } }); return container; } /** * 接受消息的監(jiān)聽(tīng),這個(gè)監(jiān)聽(tīng)會(huì)接受消息隊(duì)列1的消息 * 針對(duì)消費(fèi)者配置 * @return */ @Bean public SimpleMessageListenerContainer messageContainer2() { SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory()); container.setQueues(queue1()); container.setExposeListenerChannel(true); container.setMaxConcurrentConsumers(1); container.setConcurrentConsumers(1); container.setAcknowledgeMode(AcknowledgeMode.MANUAL); //設(shè)置確認(rèn)模式手工確認(rèn) container.setMessageListener(new ChannelAwareMessageListener() { public void onMessage(Message message, com.rabbitmq.client.Channel channel) throws Exception { byte[] body = message.getBody(); System.out.println("queue1 收到消息 : " + new String(body)); channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); //確認(rèn)消息成功消費(fèi) } }); return container; } }
注意,為了更好的展示如何配置,我配置了2個(gè)消息隊(duì)列,而本類(lèi)除了鏈接配置哪里,其他都是針對(duì)消息消費(fèi)者的,當(dāng)然不管消息消費(fèi)者和消息生產(chǎn)者都需要配置鏈接信息,而為了方便,所以本項(xiàng)目的消息消費(fèi)者和生產(chǎn)者都在本項(xiàng)目,一般實(shí)際項(xiàng)目中不會(huì)在同一項(xiàng)目,由于注釋很詳細(xì),我就不多說(shuō)了
發(fā)送消息
為了方便發(fā)送消息,所以我直接寫(xiě)了一個(gè)Controller,通過(guò)訪問(wèn)接口的形式來(lái)調(diào)用發(fā)送消息的方法,話不多說(shuō),上代碼
package wang.raye.rabbitmq.demo1; import java.util.UUID; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.support.CorrelationData; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; /** * 測(cè)試RabbitMQ發(fā)送消息的Controller * @author Raye * */ @RestController public class SendController implements RabbitTemplate.ConfirmCallback{ private RabbitTemplate rabbitTemplate; /** * 配置發(fā)送消息的rabbitTemplate,因?yàn)槭菢?gòu)造方法,所以不用注解Spring也會(huì)自動(dòng)注入(應(yīng)該是新版本的特性) * @param rabbitTemplate */ public SendController(RabbitTemplate rabbitTemplate){ this.rabbitTemplate = rabbitTemplate; //設(shè)置消費(fèi)回調(diào) this.rabbitTemplate.setConfirmCallback(this); } /** * 向消息隊(duì)列1中發(fā)送消息 * @param msg * @return */ @RequestMapping("send1") public String send1(String msg){ String uuid = UUID.randomUUID().toString(); CorrelationData correlationId = new CorrelationData(uuid); rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE, RabbitMQConfig.ROUTINGKEY1, msg, correlationId); return null; } /** * 向消息隊(duì)列2中發(fā)送消息 * @param msg * @return */ @RequestMapping("send2") public String send2(String msg){ String uuid = UUID.randomUUID().toString(); CorrelationData correlationId = new CorrelationData(uuid); rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE, RabbitMQConfig.ROUTINGKEY2, msg, correlationId); return null; } /** * 消息的回調(diào),主要是實(shí)現(xiàn)RabbitTemplate.ConfirmCallback接口 * 注意,消息回調(diào)只能代表成功消息發(fā)送到RabbitMQ服務(wù)器,不能代表消息被成功處理和接受 */ public void confirm(CorrelationData correlationData, boolean ack, String cause) { System.out.println(" 回調(diào)id:" + correlationData); if (ack) { System.out.println("消息成功消費(fèi)"); } else { System.out.println("消息消費(fèi)失敗:" + cause+"\n重新發(fā)送"); } } }
需要注意的是消息回調(diào)只能代表消息成功發(fā)送到RabbitMQ服務(wù)器
然后我們啟動(dòng)項(xiàng)目,訪問(wèn)http://localhost:8082/send1?msg=aaaa 會(huì)發(fā)現(xiàn)控制臺(tái)輸出了
收到消息 : aaaa
回調(diào)id:CorrelationData [id=37e6e913-835a-4eca-98d1-807325c5900f]
消息成功消費(fèi)
當(dāng)然回調(diào)id可能不同,如果我們?cè)L問(wèn)http://localhost:8082/send2?msg=bbbb 則輸出
queue1 收到消息 : bbbb
回調(diào)id:CorrelationData [id=0cec7500-3117-4aa2-9ea5-4790879812d4]
消息成功消費(fèi)
最后說(shuō)兩句
因?yàn)楸疚闹饕钦f(shuō)明如何從零到springboot集成RabbitMQ,所以對(duì)于RabbitMQ的很多信息和用法沒(méi)有說(shuō)明,如果對(duì)RabbitMQ本身不太熟悉的可以去看看其他關(guān)于RabbitMQ的文章,附上本文demo
以上就是本文的全部?jī)?nèi)容,希望對(duì)大家的學(xué)習(xí)有所幫助,也希望大家多多支持億速云。
免責(zé)聲明:本站發(fā)布的內(nèi)容(圖片、視頻和文字)以原創(chuàng)、轉(zhuǎn)載和分享為主,文章觀點(diǎn)不代表本網(wǎng)站立場(chǎng),如果涉及侵權(quán)請(qǐng)聯(lián)系站長(zhǎng)郵箱:is@yisu.com進(jìn)行舉報(bào),并提供相關(guān)證據(jù),一經(jīng)查實(shí),將立刻刪除涉嫌侵權(quán)內(nèi)容。