RabbitMQ整合springboot
上节我们讲了rabbitmq的概念原理5种模式已经分别与Java进行整合这节我们讲rabbitmq与springboot进行整合。要提前准备好一个spring项目的基本框架整合时有两种方式创建交换机队列及他们的绑定关系1.添加配置类加上Configuration注解和Bean注解(topic模式direct模式使用也需要用到RabbitListener只起监听队列的作用)2.用注解创建RabbitListener注解中的内容比之增加用来创建交换机队列及他们的绑定关系以及监听队列fanout模式1.引入依赖dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-amqp/artifactId/dependency然后在spring项目中创建两个模块一个生产者模块一个消费者模块2.在yml/properties文件进行相应的配置分别在两个模块中都进行创建server:port:8080spring:application:name:RabbitMQ-demo #rabbitmq配置 rabbitmq:username:guest password:guest virtual-host:/host:127.0.0.1port:56723.创建配置类创建相应的交换机与队列且绑定关系分别在两个模块中都进行创建防止一方没有配置类启动时报错packagecom.lx.producer.config;importorg.springframework.amqp.core.Binding;importorg.springframework.amqp.core.BindingBuilder;importorg.springframework.amqp.core.FanoutExchange;importorg.springframework.amqp.core.Queue;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;ConfigurationpublicclassRabbitMqConfiguration{//创建一个fanout类型的交换机BeanpublicFanoutExchangefanoutExchange(){//第一个参数是交换机的名字第二个是是否持久化第三个是是否自动删除returnnewFanoutExchange(fanout_order_exchange,true,false);}//创建三个队列BeanpublicQueuesmsQueue(){returnnewQueue(sms.fanout.queue,true);}BeanpublicQueueduanxinQueue(){returnnewQueue(duanxin.fanout.queue,true);}BeanpublicQueueemailQueue(){returnnewQueue(email.fanout.queue,true);}//三个队列分别进行绑定关系BeanpublicBindingsmsBinding(){returnBindingBuilder.bind(smsQueue()).to(fanoutExchange());}BeanpublicBindingduanxinBinding(){returnBindingBuilder.bind(duanxinQueue()).to(fanoutExchange());}BeanpublicBindingemailBinding(){returnBindingBuilder.bind(emailQueue()).to(fanoutExchange());}}4.编写业务代码在生产者模块中的service层中创建相应的生产信息业务类packagecom.lx.producer.service;importjakarta.annotation.Resource;importorg.springframework.amqp.rabbit.core.RabbitTemplate;importorg.springframework.stereotype.Service;importjava.sql.SQLOutput;importjava.util.UUID;ServicepublicclassOrderService{ResourceRabbitTemplaterabbitTemplate;publicvoidmakeOrder(StringuserId,StringproductId,Integernumbers){StringorderIdUUID.randomUUID().toString();System.out.println(订单生产成功orderId);//通过MQ完成消息的发送StringexchangeNamefanout_order_exchange;StringroutingKey;rabbitTemplate.convertAndSend(exchangeName,routingKey,orderId);}}在消费者模块中创建3个相应的消费信息业务类注意必须要在类上写上RabbitListener注解表示监听某个队列packagecom.lx.consumer.service;importorg.springframework.amqp.rabbit.annotation.RabbitHandler;importorg.springframework.amqp.rabbit.annotation.RabbitListener;importorg.springframework.stereotype.Service;ServiceRabbitListener(queues{duanxin.fanout.queue})publicclassFanoutDuanxinConsumer{RabbitHandlerpublicvoidreceiveMessage(Stringmessage){System.out.println(duanxin.fanout---接收到的信息是message);}}packagecom.lx.consumer.service;importorg.springframework.amqp.rabbit.annotation.RabbitHandler;importorg.springframework.amqp.rabbit.annotation.RabbitListener;importorg.springframework.stereotype.Service;ServiceRabbitListener(queues{email.fanout.queue})publicclassFanoutEmailConsumer{RabbitHandlerpublicvoidreceiveMessage(Stringmessage){System.out.println(email.fanout---接收到的信息是message);}}packagecom.lx.consumer.service;importorg.springframework.amqp.rabbit.annotation.RabbitHandler;importorg.springframework.amqp.rabbit.annotation.RabbitListener;importorg.springframework.stereotype.Service;ServiceRabbitListener(queues{sms.fanout.queue})publicclassFanoutSmsConsumer{RabbitHandlerpublicvoidreceiveMessage(Stringmessage){System.out.println(sms.fanout---接收到的信息是message);}}最后先启动rabbitmq在启动两个项目我们就能得到再消费者中得到生产者发送的消息了direct模式direct模式与fanout模式的代码差不多只是要把配置类和生产者的业务处理类修改一下配置类修改如下把刚刚fanout命名都改为direct了packagecom.lx.producer.config;importorg.springframework.amqp.core.*;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;ConfigurationpublicclassRabbitMqConfiguration{//创建一个direct类型的交换机BeanpublicDirectExchangedirectExchangeExchange(){//第一个参数是交换机的名字第二个是是否持久化第三个是是否自动删除returnnewDirectExchange(direct_order_exchange,true,false);}//创建三个队列BeanpublicQueuesmsQueue(){returnnewQueue(sms.direct.queue,true);}BeanpublicQueueduanxinQueue(){returnnewQueue(duanxin.direct.queue,true);}BeanpublicQueueemailQueue(){returnnewQueue(email.direct.queue,true);}//三个队列分别进行绑定关系BeanpublicBindingsmsBinding(){returnBindingBuilder.bind(smsQueue()).to(directExchangeExchange()).with(sms);}BeanpublicBindingduanxinBinding(){returnBindingBuilder.bind(duanxinQueue()).to(directExchangeExchange()).with(duanxin);}BeanpublicBindingemailBinding(){returnBindingBuilder.bind(emailQueue()).to(directExchangeExchange()).with(email);}}改成穿件direct类型的交换机以及路由与队列绑定时加了key生产者业务类packagecom.lx.producer.service;importjakarta.annotation.Resource;importorg.springframework.amqp.rabbit.core.RabbitTemplate;importorg.springframework.stereotype.Service;importjava.sql.SQLOutput;importjava.util.UUID;ServicepublicclassOrderService{ResourceRabbitTemplaterabbitTemplate;publicvoidmakeOrder(StringuserId,StringproductId,Integernumbers){StringorderIdUUID.randomUUID().toString();System.out.println(订单生产成功orderId);//通过MQ完成消息的发送StringexchangeNamedirect_order_exchange;StringroutingKey1sms;StringroutingKey2duanxin;rabbitTemplate.convertAndSend(exchangeName,routingKey1,orderId);rabbitTemplate.convertAndSend(exchangeName,routingKey2,orderId);}}其余的代码不要变把所有名字从fanout改为direct就行topic模式topic模式中我们使用注解的方式来创建交换机队列以及它们之间的绑定关系只需要把消费者业务代码及生产者业务代码修改下即可消费者业务代码(修改RabbitListener注解)packagecom.lx.consumer.service;importorg.springframework.amqp.core.ExchangeTypes;importorg.springframework.amqp.rabbit.annotation.*;importorg.springframework.stereotype.Service;ServiceRabbitListener(bindingsQueueBinding(valueQueue(valueduanxin.topic.queue,durabletrue,autoDeletefalse),exchangeExchange(valuetopic_order_exchange,typeExchangeTypes.TOPIC),key#.duanxin.#))publicclassTopicDuanxinConsumer{RabbitHandlerpublicvoidreceiveMessage(Stringmessage){System.out.println(duanxin.topic---接收到的信息是message);}}packagecom.lx.consumer.service;importorg.springframework.amqp.core.ExchangeTypes;importorg.springframework.amqp.rabbit.annotation.*;importorg.springframework.stereotype.Service;ServiceRabbitListener(bindingsQueueBinding(valueQueue(valuetopic.email.queue,durabletrue,autoDeletefalse),exchangeExchange(valuetopic_order_exchange,typeExchangeTypes.TOPIC),key*.email.#))publicclassTopicEmailConsumer{RabbitHandlerpublicvoidreceiveMessage(Stringmessage){System.out.println(email.direct---接收到的信息是message);}}packagecom.lx.consumer.service;importorg.springframework.amqp.core.ExchangeTypes;importorg.springframework.amqp.rabbit.annotation.*;importorg.springframework.stereotype.Service;ServiceRabbitListener(bindingsQueueBinding(valueQueue(valuesms.topic.queue,durabletrue,autoDeletefalse),exchangeExchange(valuetopic_order_exchange,typeExchangeTypes.TOPIC),keycom.#))publicclassTopicSmsConsumer{RabbitHandlerpublicvoidreceiveMessage(Stringmessage){System.out.println(sms.direct---接收到的信息是message);}}生产者业务代码(只修改routing keyj就行通配符模式)packagecom.lx.producer.service;importjakarta.annotation.Resource;importorg.springframework.amqp.rabbit.core.RabbitTemplate;importorg.springframework.stereotype.Service;importjava.sql.SQLOutput;importjava.util.UUID;ServicepublicclassOrderService{ResourceRabbitTemplaterabbitTemplate;publicvoidmakeOrder(StringuserId,StringproductId,Integernumbers){StringorderIdUUID.randomUUID().toString();System.out.println(订单生产成功orderId);//通过MQ完成消息的发送StringexchangeNametopic_order_exchange;StringroutingKey1com.email.duanxin;// String routingKey2duanxin;rabbitTemplate.convertAndSend(exchangeName,routingKey1,orderId);// rabbitTemplate.convertAndSend(exchangeName,routingKey2,orderId);}}