RabbitMQ簡介

創(chuàng)新互聯(lián)公司主要為客戶提供服務(wù)項目涵蓋了網(wǎng)頁視覺設(shè)計、VI標(biāo)志設(shè)計、全網(wǎng)營銷推廣、網(wǎng)站程序開發(fā)、HTML5響應(yīng)式重慶網(wǎng)站建設(shè)公司、移動網(wǎng)站建設(shè)、微商城、網(wǎng)站托管及成都網(wǎng)站維護(hù)、WEB系統(tǒng)開發(fā)、域名注冊、國內(nèi)外服務(wù)器租用、視頻、平面設(shè)計、SEO優(yōu)化排名。設(shè)計、前端、后端三個建站步驟的完善服務(wù)體系。一人跟蹤測試的建站服務(wù)標(biāo)準(zhǔn)。已經(jīng)為葡萄架行業(yè)客戶提供了網(wǎng)站建設(shè)服務(wù)。
RabbitMQ是一個在AMQP基礎(chǔ)上完整的,可復(fù)用的企業(yè)消息系統(tǒng)
MQ全稱為Message Queue, 消息隊列(MQ)是一種應(yīng)用程序?qū)?yīng)用程序的通信方法。應(yīng)用程序通過讀寫出入隊列的消息(針對應(yīng)用程序的數(shù)據(jù))來通信,而無需專用連接來鏈接它們。消息傳遞指的是程序之間通過在消息中發(fā)送數(shù)據(jù)進(jìn)行通信,而不是通過直接調(diào)用彼此來通信,直接調(diào)用通常是用于諸如遠(yuǎn)程過程調(diào)用的技術(shù)。排隊指的是應(yīng)用程序通過 隊列來通信。隊列的使用除去了接收和發(fā)送應(yīng)用程序同時執(zhí)行的要求。
AMQP就是一個協(xié)議,是一個高級抽象層消息通信協(xié)議。
雖然在同步消息通訊的世界里有很多公開標(biāo)準(zhǔn)(如 COBAR的 IIOP ,或者是 SOAP 等),但是在異步消息處理中卻不是這樣,只有大企業(yè)有一些商業(yè)實現(xiàn)(如微軟的 MSMQ ,IBM 的 Websphere MQ 等),因此,在 2006 年的 6 月,Cisco 、Redhat、iMatix 等聯(lián)合制定了 AMQP 的公開標(biāo)準(zhǔn)。也就是說AMQP是異步通訊的一個協(xié)議。
RabbitMQ使用場景
在項目中,將一些無需即時返回且耗時的操作提取出來,進(jìn)行了異步處理,而這種異步處理的方式大大的節(jié)省了服務(wù)器的請求響應(yīng)時間,從而提高了系統(tǒng)的吞吐量。不過大多數(shù)不僅僅是無需即時返回,甚至是執(zhí)行是否成功都無所謂。如果需要即時返回則可以使用Dubbo,Spring boot與Dubbo集成可以去看Spring boot 集成Dubbox
RabbitMQ依賴
RabbitMQ并不是直接一個簡單的jar包(Jar包只是提供一個基本的與RabbitMQ本身通訊的一些功能),和Dubbo相同,RabbitMQ也需要其他軟件來運行,以下是RabbitMQ運行所需要的軟件
1、Erlang
由于RabbitMQ軟件本身是基于Erlang開發(fā)的,所以想要運行RabbitMQ必須要先按照Erlang
Erlang官網(wǎng)
Erlang下載地址
RabbitMQ
RabbitMQ才是實現(xiàn)消息隊列的核心
RabbitMQ官網(wǎng)
RabbitMQ下載
配置RabbitMQ
安裝完成后,需要完成一些配置才能使用RabbitMQ,可以直接用cmd到RabbitMQ的安裝目錄下的sbin目錄通過命令配置,也可以直接在開始菜單中直接找到RabbitMQ Command Prompt (sbin dir)運行直接到達(dá)RabbitMQ的安裝目錄的sbin,為了方便,我們先啟用管理插件,執(zhí)行命令
rabbitmq-plugins.bat enable rabbitmq_management
即可,注意,這是在Windows下面,如果是Linux則沒有bat后綴然后我們添加一個用戶,因為在外網(wǎng)環(huán)境沒有用戶的情況下是不能連接成功的,執(zhí)行添加用戶命令
rabbitmqctl.bat add_user springboot password
springboot是用戶名,password是密碼
然后為了方便演示,我們給springboot賦予管理員權(quán)限,方便登錄管理頁面
rabbitmqctl.bat set_user_tags springboot administrator
給賬號賦予虛擬主機(jī)權(quán)限
rabbitmqctl.bat set_permissions -p / springboot .* .* .*
然后啟動RabbitMQ服務(wù) 訪問RabbitMQ管理頁面http://localhost:15672即可看見登錄頁面,如果沒有創(chuàng)建用戶則可以用guest,guest登錄,如果有創(chuàng)建用戶則用創(chuàng)建的用戶登錄
創(chuàng)建Springboot項目
因為創(chuàng)建spring boot項目在前面的文章已經(jīng)說過很多次了,所以這里就不多說了
添加RabbitMQ相關(guān)依賴
<!-- rabbitmq -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>沒錯,就是點配置,不過這樣可能有點不理解,我還是把全部配置貼出來吧
<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>
因為沒有做其他操作,所以目前項目主要是依賴2個模塊,一個Sprig boot,一個RabbitMQ
添加配置類
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 的配置類
*
* @author Raye
* @since 2016年10月12日10:57:44
*/
@Configuration
public class RabbitMQConfig {
/** 消息交換機(jī)的名字*/
public static final String EXCHANGE = "my-mq-exchange";
/** 隊列key1*/
public static final String ROUTINGKEY1 = "queue_one_key1";
/** 隊列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ī)
* 針對消費者配置
FanoutExchange: 將消息分發(fā)到所有的綁定隊列,無routingkey的概念
HeadersExchange :通過添加屬性key-value匹配
DirectExchange:按照routingkey分發(fā)到指定隊列
TopicExchange:多關(guān)鍵字匹配
*/
@Bean
public DirectExchange defaultExchange() {
return new DirectExchange(EXCHANGE, true, false);
}
/**
* 配置消息隊列1
* 針對消費者配置
* @return
*/
@Bean
public Queue queue() {
return new Queue("queue_one", true); //隊列持久
}
/**
* 將消息隊列1與交換機(jī)綁定
* 針對消費者配置
* @return
*/
@Bean
public Binding binding() {
return BindingBuilder.bind(queue()).to(defaultExchange()).with(RabbitMQConfig.ROUTINGKEY1);
}
/**
* 配置消息隊列2
* 針對消費者配置
* @return
*/
@Bean
public Queue queue1() {
return new Queue("queue_one1", true); //隊列持久
}
/**
* 將消息隊列2與交換機(jī)綁定
* 針對消費者配置
* @return
*/
@Bean
public Binding binding1() {
return BindingBuilder.bind(queue1()).to(defaultExchange()).with(RabbitMQConfig.ROUTINGKEY2);
}
/**
* 接受消息的監(jiān)聽,這個監(jiān)聽會接受消息隊列1的消息
* 針對消費者配置
* @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)消息成功消費
}
});
return container;
}
/**
* 接受消息的監(jiān)聽,這個監(jiān)聽會接受消息隊列1的消息
* 針對消費者配置
* @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)消息成功消費
}
});
return container;
}
}
注意,為了更好的展示如何配置,我配置了2個消息隊列,而本類除了鏈接配置哪里,其他都是針對消息消費者的,當(dāng)然不管消息消費者和消息生產(chǎn)者都需要配置鏈接信息,而為了方便,所以本項目的消息消費者和生產(chǎn)者都在本項目,一般實際項目中不會在同一項目,由于注釋很詳細(xì),我就不多說了
發(fā)送消息
為了方便發(fā)送消息,所以我直接寫了一個Controller,通過訪問接口的形式來調(diào)用發(fā)送消息的方法,話不多說,上代碼
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;
/**
* 測試RabbitMQ發(fā)送消息的Controller
* @author Raye
*
*/
@RestController
public class SendController implements RabbitTemplate.ConfirmCallback{
private RabbitTemplate rabbitTemplate;
/**
* 配置發(fā)送消息的rabbitTemplate,因為是構(gòu)造方法,所以不用注解Spring也會自動注入(應(yīng)該是新版本的特性)
* @param rabbitTemplate
*/
public SendController(RabbitTemplate rabbitTemplate){
this.rabbitTemplate = rabbitTemplate;
//設(shè)置消費回調(diào)
this.rabbitTemplate.setConfirmCallback(this);
}
/**
* 向消息隊列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;
}
/**
* 向消息隊列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),主要是實現(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("消息成功消費");
} else {
System.out.println("消息消費失敗:" + cause+"\n重新發(fā)送");
}
}
}
需要注意的是消息回調(diào)只能代表消息成功發(fā)送到RabbitMQ服務(wù)器
然后我們啟動項目,訪問http://localhost:8082/send1?msg=aaaa 會發(fā)現(xiàn)控制臺輸出了
收到消息 : aaaa
回調(diào)id:CorrelationData [id=37e6e913-835a-4eca-98d1-807325c5900f]
消息成功消費
當(dāng)然回調(diào)id可能不同,如果我們訪問http://localhost:8082/send2?msg=bbbb 則輸出
queue1 收到消息 : bbbb
回調(diào)id:CorrelationData [id=0cec7500-3117-4aa2-9ea5-4790879812d4]
消息成功消費
最后說兩句
因為本文主要是說明如何從零到springboot集成RabbitMQ,所以對于RabbitMQ的很多信息和用法沒有說明,如果對RabbitMQ本身不太熟悉的可以去看看其他關(guān)于RabbitMQ的文章,附上本文demo
以上就是本文的全部內(nèi)容,希望對大家的學(xué)習(xí)有所幫助,也希望大家多多支持創(chuàng)新互聯(lián)。
文章標(biāo)題:Springboot集成RabbitMQ的示例代碼
瀏覽地址:http://chinadenli.net/article0/gdjjio.html
成都網(wǎng)站建設(shè)公司_創(chuàng)新互聯(lián),為您提供商城網(wǎng)站、網(wǎng)頁設(shè)計公司、靜態(tài)網(wǎng)站、微信小程序、定制開發(fā)、移動網(wǎng)站建設(shè)
聲明:本網(wǎng)站發(fā)布的內(nèi)容(圖片、視頻和文字)以用戶投稿、用戶轉(zhuǎn)載內(nèi)容為主,如果涉及侵權(quán)請盡快告知,我們將會在第一時間刪除。文章觀點不代表本網(wǎng)站立場,如需處理請聯(lián)系客服。電話:028-86922220;郵箱:631063699@qq.com。內(nèi)容未經(jīng)允許不得轉(zhuǎn)載,或轉(zhuǎn)載時需注明來源: 創(chuàng)新互聯(lián)