最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Spring boot集成RabbitMQ的示例代碼

 更新時間:2018年05月20日 10:15:24   作者:Raye Blog  
本篇文章主要介紹了Spring boot集成RabbitMQ的示例代碼,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧

RabbitMQ簡介

RabbitMQ是一個在AMQP基礎上完整的,可復用的企業(yè)消息系統(tǒng)

MQ全稱為Message Queue, 消息隊列(MQ)是一種應用程序?qū)贸绦虻耐ㄐ欧椒?。應用程序通過讀寫出入隊列的消息(針對應用程序的數(shù)據(jù))來通信,而無需專用連接來鏈接它們。消息傳遞指的是程序之間通過在消息中發(fā)送數(shù)據(jù)進行通信,而不是通過直接調(diào)用彼此來通信,直接調(diào)用通常是用于諸如遠程過程調(diào)用的技術(shù)。排隊指的是應用程序通過 隊列來通信。隊列的使用除去了接收和發(fā)送應用程序同時執(zhí)行的要求。

AMQP就是一個協(xié)議,是一個高級抽象層消息通信協(xié)議。

雖然在同步消息通訊的世界里有很多公開標準(如 COBAR的 IIOP ,或者是 SOAP 等),但是在異步消息處理中卻不是這樣,只有大企業(yè)有一些商業(yè)實現(xiàn)(如微軟的 MSMQ ,IBM 的 Websphere MQ 等),因此,在 2006 年的 6 月,Cisco 、Redhat、iMatix 等聯(lián)合制定了 AMQP 的公開標準。也就是說AMQP是異步通訊的一個協(xié)議。

RabbitMQ使用場景

在項目中,將一些無需即時返回且耗時的操作提取出來,進行了異步處理,而這種異步處理的方式大大的節(jié)省了服務器的請求響應時間,從而提高了系統(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)運行直接到達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 

給賬號賦予虛擬主機權(quán)限

rabbitmqctl.bat set_permissions -p / springboot .* .* .* 

然后啟動RabbitMQ服務 訪問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 { 
  /** 消息交換機的名字*/
  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); // 必須要設置
    return connectionFactory;
  }

  /** 
   * 配置消息交換機
   * 針對消費者配置 
    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與交換機綁定
   * 針對消費者配置 
   * @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與交換機綁定
   * 針對消費者配置 
   * @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); //設置確認模式手工確認 
    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); //確認消息成功消費 

      } 

    }); 
    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); //設置確認模式手工確認 
    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); //確認消息成功消費 
      } 
    }); 
    return container; 
  } 
}

注意,為了更好的展示如何配置,我配置了2個消息隊列,而本類除了鏈接配置哪里,其他都是針對消息消費者的,當然不管消息消費者和消息生產(chǎn)者都需要配置鏈接信息,而為了方便,所以本項目的消息消費者和生產(chǎn)者都在本項目,一般實際項目中不會在同一項目,由于注釋很詳細,我就不多說了

發(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也會自動注入(應該是新版本的特性)
   * @param rabbitTemplate
   */
  public SendController(RabbitTemplate rabbitTemplate){
    this.rabbitTemplate = rabbitTemplate;
    //設置消費回調(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服務器,不能代表消息被成功處理和接受
   */
  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服務器

然后我們啟動項目,訪問http://localhost:8082/send1?msg=aaaa 會發(fā)現(xiàn)控制臺輸出了

收到消息 : aaaa
 回調(diào)id:CorrelationData [id=37e6e913-835a-4eca-98d1-807325c5900f]
消息成功消費

當然回調(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)容,希望對大家的學習有所幫助,也希望大家多多支持腳本之家。

  • JAVA常用API總結(jié)與說明

    JAVA常用API總結(jié)與說明

    這篇文章主要介紹了JAVA常用API總結(jié)與說明,包括JAVA線程常用API,JAVA隊列常用API,JAVA泛型集合算法常用API,JAVA并發(fā)常用API需要的朋友可以參考下
    2022-12-12
  • java處理csv文件上傳示例詳解

    java處理csv文件上傳示例詳解

    這篇文章主要為大家詳細介紹了java處理csv文件上傳示例,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2020-12-12
  • SpringBoot整合RabbitMQ, 實現(xiàn)生產(chǎn)者與消費者的功能

    SpringBoot整合RabbitMQ, 實現(xiàn)生產(chǎn)者與消費者的功能

    這篇文章主要介紹了SpringBoot整合RabbitMQ, 實現(xiàn)生產(chǎn)者與消費者的功能,幫助大家更好得理解和學習使用SpringBoot框架,感興趣的朋友可以了解下
    2021-03-03
  • SPFA算法的實現(xiàn)原理及其應用詳解

    SPFA算法的實現(xiàn)原理及其應用詳解

    SPFA算法,全稱為Shortest?Path?Faster?Algorithm,是求解單源最短路徑問題的一種常用算法,本文就來聊聊它的實現(xiàn)原理與簡單應用吧
    2023-05-05
  • Spring事務失效場景的詳細整理

    Spring事務失效場景的詳細整理

    Spring 事務的傳播特性說的是,當多個事務同時存在的時候,Spring 如何處理這些事務的特性,下面這篇文章主要給大家介紹了關(guān)于Spring事務失效場景的相關(guān)資料,需要的朋友可以參考下
    2022-02-02
  • mybatis中orderBy(排序字段)和sort(排序方式)引起的bug及解決

    mybatis中orderBy(排序字段)和sort(排序方式)引起的bug及解決

    這篇文章主要介紹了mybatis中orderBy(排序字段)和sort(排序方式)引起的bug,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-01-01
  • 詳解設計模式在Spring中的應用(9種)

    詳解設計模式在Spring中的應用(9種)

    這篇文章主要介紹了詳解設計模式在Spring中的應用(9種),詳細的介紹了這9種模式在項目中的應用,具有一定的參考價值,感興趣的可以了解一下
    2019-04-04
  • Java 實現(xiàn)多線程的幾種方式匯總

    Java 實現(xiàn)多線程的幾種方式匯總

    JAVA多線程實現(xiàn)方式主要有三種:繼承Thread類、實現(xiàn)Runnable接口、使用ExecutorService、Callable、Future實現(xiàn)有返回結(jié)果的多線程。其中前兩種方式線程執(zhí)行完后都沒有返回值,只有最后一種是帶返回值的。
    2016-03-03
  • 5個JAVA入門必看的經(jīng)典實例

    5個JAVA入門必看的經(jīng)典實例

    這篇文章主要為大家詳細介紹了5個JAVA入門必看的經(jīng)典實例,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2017-10-10
  • 最新評論

    务川| 临泉县| 嘉鱼县| 福海县| 公主岭市| 吴江市| 呼图壁县| 葵青区| 建昌县| 独山县| 新乡县| 巴青县| 兴国县| 清镇市| 饶平县| 盘山县| 鹤岗市| 丰县| 宁陵县| 板桥市| 建昌县| 天长市| 邓州市| 酒泉市| 五常市| 乐平市| 拜泉县| 吉林省| 宁阳县| 泸溪县| 乳源| 沙洋县| 辽源市| 泽库县| 罗源县| 宿松县| 涪陵区| 韩城市| 石门县| 扎鲁特旗| 晴隆县|