springboot2 springcloud Greenwich.SR3 构建微服务--5.rabbitMQ的使用

本文主要是介绍springboot2 springcloud Greenwich.SR3 构建微服务--5.rabbitMQ的使用,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

其实在统一配置中心的时候就应该开始说rabbitmq 了, 但是没有说, 所以这里专门说下rabbitmq.

rabbitmq 在配置中心作为消息的传递来使用了.

 

rabbitmq的原理, 具体使用, 请查阅这个

https://blog.csdn.net/u010398771/article/details/84136959

 

现在的mq开源的不要太多了, 我最先用的activemq, 后来阿里开源的rocketmq,

还有 kafka, rabbitmq. 最近看新闻,腾讯也开源了他们用java写的MQ, 名字不记得了...又是个亿万级别的服务, 牛掰哄哄的样子

其实rabbitMQ的安装很简单, 这里就不推荐大家先下载erlang 再安装rabbitmq 了, 更不推荐你去编译源码安装, 因为我安装过一次, 装的时候,电脑卡, 时间又长, 你可以直接在你的docker里面安装一个rabbitmq, 就可以使用了,很方便.

docker安装rabbitmq, 先搜索目前的镜像

docker search rabbitmq

先拉取下来, 再创建容器运行.注意拉取带有management 的rabbitmq , 这是带控制带的版本.

#拉取镜像
docker pull rabbitmq:3.7.17-management#运行
docker run -d -p 5672:5672 -p 15672:15672 --name myrabbitmq rabbitmq:3.7.17-management

登陆网页端, http://localhost:15672,  用户名和密码都是guest, 就可以看到了管理页面了, 可以看到目前没有消息到来和消费

 

rabbitmq的使用

1.添加依赖

 <dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

2.添加配置文件

#rabbitmq的配置
spring.rabbitmq.host=127.0.0.1
spring.rabbitmq.port=5672
spring.rabbitmq.username=guest
spring.rabbitmq.password=guest
#解决input 和 output中重名bean问题,需要配置这个....
spring.cloud.stream.bindings.input.destination: raw-sensor-data
spring.cloud.stream.bindings.output.destination: raw-sensor-data

 

2.1. queue 的  bean的配置

//名字都变成常量
public class MessageConstant {public final static String INPUT = "input";public final static String OUTPUT = "output";public final static String QUEUE_NAME = "myQueue";}
/*** 注入一个queue, 不要去手动创建!*/@Beanpublic Queue queue(){return new Queue(MessageConstant.QUEUE_NAME,true);}

3.就开始发送和接收消息了.


import com.example.order.constans.MessageConstant;
import com.example.order.dto.OrderDTO;
import com.example.order.message.StreamClient;
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;import java.time.LocalDateTime;/*** @author zk* @Description:* @date 2019-09-20 16:24*/
@RestController
public class SendMessageController {/*** 往rabbitmq发送消息*/@Autowiredprivate AmqpTemplate amqpTemplate;@GetMapping("send")public String sendMessage() {//需要实现创建一个myQueue的队列amqpTemplate.convertAndSend(MessageConstant.QUEUE_NAME, "this is the first message");
//第二种消息发送amqpTemplate.convertAndSend("myOrder", "computer", "this is the first computer  message");amqpTemplate.convertAndSend("myOrder", "fruit", "this is the first message");return "success";}@Autowired(required = false)private StreamClient streamClient;/*** 发送字符串对象** @return*//*@GetMapping("send2")public String send2() {String msg = "time  is today, haha ";System.out.println("发送的消息是:" + msg);streamClient.output().send(MessageBuilder.withPayload(msg).build());return "success";}*//*** 发送java bean对象* 这个send3 不能和send2 使用的同一个变量.但是不能共存, 需要注释掉, 不然报错的...** @return*/@GetMapping("send3")public String send3(){OrderDTO dto = new OrderDTO();dto.setBuyerAddress("address");dto.setBuyerName("xiexin");streamClient.output().send(MessageBuilder.withPayload(dto).build());return "success";}}

4.消息接收端


import com.example.order.constans.MessageConstant;
import com.example.order.dto.OrderDTO;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.messaging.handler.annotation.SendTo;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;/*** @author zk* @Description: 接受消息* @date 2019-09-20 16:21*/
@Component
@EnableBinding(StreamClient.class)
public class StreamReceiver {/*** 接受字符串*//*@StreamListener(MessageConstant.INPUT)public void process(Object msg) {System.out.println(msg);}*//*** 接受orderDTO对象*/@StreamListener(MessageConstant.INPUT)@SendTo(MessageConstant.OUTPUT)//处理完消息,再回发送个消息public Object process(OrderDTO msg) {System.out.println(msg);//发送消息return "这是消息";}/*** 接受上面的回传的消息*/@StreamListener(MessageConstant.OUTPUT)public void process2(String msg) {System.out.println(msg);}}
import com.example.order.constans.MessageConstant;
import org.springframework.amqp.rabbit.annotation.Exchange;
import org.springframework.amqp.rabbit.annotation.Queue;
import org.springframework.amqp.rabbit.annotation.QueueBinding;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;/*** @author zk* @Description:mq 接收消息* @date 2019-09-20 15:29*/
@Component
public class MQReceiver {//1.手动创建队列//@RabbitListener(queues= MessageConstant.QUEUE_NAME)//2.自动创建队列//@RabbitListener(queuesToDeclare = @Queue(MessageConstant.QUEUE_NAME))//3.自动创建 exchange 和 queue绑定@RabbitListener(bindings = @QueueBinding(value = @Queue(MessageConstant.QUEUE_NAME),exchange = @Exchange("myExchange")))public void process(String message){System.out.println("接收到的消息是:"+message);}/*** 按照exchange 和 key 进行匹配.* 数码消息接受商* @param message*/@RabbitListener(bindings = @QueueBinding(exchange = @Exchange("myOrder"),key = "computer",value = @Queue("computerOrder")))public void computerProcess(String message){System.out.println("computer接收到的消息是:"+message);}/*** 水果消息接受商* @param message*/@RabbitListener(bindings = @QueueBinding(exchange = @Exchange("myOrder"),key = "fruit",value = @Queue("fruitOrder")))public void fruitProcess(String message){System.out.println("fruit接收到的消息是:"+message);}}
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.TypeReference;
import com.example.product.entity.ProductInfo;
import org.springframework.amqp.rabbit.annotation.Queue;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;import java.util.HashMap;
import java.util.List;
import java.util.Map;@Component
public class ProductInfoReceiver {private final static String HASH_KEY="product_stock";@Autowiredprivate StringRedisTemplate stringRedisTemplate;/*** 传递整个对象的方法* @param message*///@RabbitListener(queuesToDeclare = @Queue("productInfo"))public void process0(String message) {//注意反序列化的时候,两个对象的前缀要求是一样的(包名)//ProductInfo info = JSON.parseObject(message, ProductInfo.class);List<ProductInfo> infos = JSON.parseObject(message, new TypeReference<List<ProductInfo>>() {});System.out.println("接受到的消息是:" + infos);//存储到redis中去.//stringRedisTemplate.opsForValue()// .set("product_stock_" + info.getProductId(), String.valueOf(info.getProductStock()));/*** 单值的key -value 可以储存, 现在多个的时候, 最好使用hash 的方式进行存储* 这种单值的存储方式不好...* product_stock        1       5*      key           field   value*/for (ProductInfo info : infos) {stringRedisTemplate.opsForHash().put(HASH_KEY,info.getProductId(),info.getProductStock());}}/*** 只传递id和剩余库存的方法* 这里的参数的接收对象一定和发送的一致,不然rabbitmq会报类型转换错误.* json转换后,就是json的数组或者json字符串了.上面的那个方法是个反例* @param message*/@RabbitListener(queuesToDeclare = @Queue("productInfo"))public void process(Map<String,Integer> message) {//注意反序列化的时候,两个对象的前缀要求是一样的(包名)System.out.println("接受到的消息是:" + message);//存储到redis中去.//stringRedisTemplate.opsForValue()// .set("product_stock_" + info.getProductId(), String.valueOf(info.getProductStock()));/*** 单值的key -value 可以储存, 现在多个的时候, 最好使用hash 的方式进行存储* 这种单值的存储方式不好...* product_stock        1       5*      key           field   value*/for (Map.Entry<String, Integer> entry : message.entrySet()) {stringRedisTemplate.opsForHash().put(HASH_KEY,entry.getKey(),String.valueOf(entry.getValue()));}}}
import com.example.order.constans.MessageConstant;
import org.springframework.cloud.stream.annotation.Input;
import org.springframework.cloud.stream.annotation.Output;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.SubscribableChannel;/*** @author zk* @Description:* @date 2019-09-20 16:19*/
public interface StreamClient {/*** 这里的input 和 output 必须是不同的字符串, 不然是同名的bean错误*/@Input(MessageConstant.INPUT)SubscribableChannel input();@Output(MessageConstant.OUTPUT)MessageChannel output();}

rabbitmq 的消息是根据 queue  exchange  key 做的分类存储,

所以上面的myqueue中的不同的key 会得到不同的接收者那里去.

 

 

 突然觉得我好懒啊, 不愿意好好写博客, 想各种偷懒, 这东西更像在写自己的随手笔记了,而不是在写让大部分都能够看懂的博客. 

随便代码往这里一贴, 然后加几句解释就完事了, 不讲原理, 因为不想打字.  不想说流程, 因为懒的去画很好看的图出来, 懒的理由还很充分! .

甚至这一系列的博文中可能一些小章节中有误,  看了代码应该运行不起来, 但是我还是贴出来, 某些还涉及到其他地方的点滴, 没有在博客中体现出来, 感觉我就是在发垃圾博客!!!! 

这次一口气写十篇关于springcloud 的博文, 但是现在看起来, 估计写完后质量非常的差. 完全不是我想象的那样子,既能够自己总结, 又能帮助别人提高.

反思下自己!!!! 

不过写文章确实很花时间, 就这种水平的博客, 想写一个出来, 也差不多得花10分钟 ! 所以有了对这10篇博客的重新修改

 

整个代码地址是:

https://github.com/changhe626/micro-service

Java Framework,欢迎各位前来交流java相关
QQ群:965125360
 

 

 

这篇关于springboot2 springcloud Greenwich.SR3 构建微服务--5.rabbitMQ的使用的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/986270

相关文章

vue使用docxtemplater导出word

《vue使用docxtemplater导出word》docxtemplater是一种邮件合并工具,以编程方式使用并处理条件、循环,并且可以扩展以插入任何内容,下面我们来看看如何使用docxtempl... 目录docxtemplatervue使用docxtemplater导出word安装常用语法 封装导出方

Linux换行符的使用方法详解

《Linux换行符的使用方法详解》本文介绍了Linux中常用的换行符LF及其在文件中的表示,展示了如何使用sed命令替换换行符,并列举了与换行符处理相关的Linux命令,通过代码讲解的非常详细,需要的... 目录简介检测文件中的换行符使用 cat -A 查看换行符使用 od -c 检查字符换行符格式转换将

Java编译生成多个.class文件的原理和作用

《Java编译生成多个.class文件的原理和作用》作为一名经验丰富的开发者,在Java项目中执行编译后,可能会发现一个.java源文件有时会产生多个.class文件,从技术实现层面详细剖析这一现象... 目录一、内部类机制与.class文件生成成员内部类(常规内部类)局部内部类(方法内部类)匿名内部类二、

SpringBoot实现数据库读写分离的3种方法小结

《SpringBoot实现数据库读写分离的3种方法小结》为了提高系统的读写性能和可用性,读写分离是一种经典的数据库架构模式,在SpringBoot应用中,有多种方式可以实现数据库读写分离,本文将介绍三... 目录一、数据库读写分离概述二、方案一:基于AbstractRoutingDataSource实现动态

Python FastAPI+Celery+RabbitMQ实现分布式图片水印处理系统

《PythonFastAPI+Celery+RabbitMQ实现分布式图片水印处理系统》这篇文章主要为大家详细介绍了PythonFastAPI如何结合Celery以及RabbitMQ实现简单的分布式... 实现思路FastAPI 服务器Celery 任务队列RabbitMQ 作为消息代理定时任务处理完整

使用Jackson进行JSON生成与解析的新手指南

《使用Jackson进行JSON生成与解析的新手指南》这篇文章主要为大家详细介绍了如何使用Jackson进行JSON生成与解析处理,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 目录1. 核心依赖2. 基础用法2.1 对象转 jsON(序列化)2.2 JSON 转对象(反序列化)3.

Springboot @Autowired和@Resource的区别解析

《Springboot@Autowired和@Resource的区别解析》@Resource是JDK提供的注解,只是Spring在实现上提供了这个注解的功能支持,本文给大家介绍Springboot@... 目录【一】定义【1】@Autowired【2】@Resource【二】区别【1】包含的属性不同【2】@

springboot循环依赖问题案例代码及解决办法

《springboot循环依赖问题案例代码及解决办法》在SpringBoot中,如果两个或多个Bean之间存在循环依赖(即BeanA依赖BeanB,而BeanB又依赖BeanA),会导致Spring的... 目录1. 什么是循环依赖?2. 循环依赖的场景案例3. 解决循环依赖的常见方法方法 1:使用 @La

Java枚举类实现Key-Value映射的多种实现方式

《Java枚举类实现Key-Value映射的多种实现方式》在Java开发中,枚举(Enum)是一种特殊的类,本文将详细介绍Java枚举类实现key-value映射的多种方式,有需要的小伙伴可以根据需要... 目录前言一、基础实现方式1.1 为枚举添加属性和构造方法二、http://www.cppcns.co

使用Python实现快速搭建本地HTTP服务器

《使用Python实现快速搭建本地HTTP服务器》:本文主要介绍如何使用Python快速搭建本地HTTP服务器,轻松实现一键HTTP文件共享,同时结合二维码技术,让访问更简单,感兴趣的小伙伴可以了... 目录1. 概述2. 快速搭建 HTTP 文件共享服务2.1 核心思路2.2 代码实现2.3 代码解读3.