赞
踩
RabbitMQ相关知识点思维导图如下:
注意一点,在发送消息的时候对template进行配置mandatory=true保证监听有效;生产端还可以配置其他属性,比如发送重试,超时时间、次数、间隔等
消费端核心配置
@RabbitListener注解的使用
消费端监听@RabbitListener注解,这个对于在实际工作中非常的好用;@RabbitListener是一个组合注解,里面可以注解配置(@QueueBinding、@Queue、@Exchange)直接通过这个组合注解一次性搞定消费端交换机、队列、绑定、路由、并且配置监听功能等
相关代码
order.java(注意一定要实现序列号接口):
- package com.ue.entity;
-
- import java.io.Serializable;
-
- public class Order implements Serializable {
-
- private String id;
- private String name;
-
- public Order() {
- }
- public Order(String id, String name) {
- super();
- this.id = id;
- this.name = name;
- }
- public String getId() {
- return id;
- }
- public void setId(String id) {
- this.id = id;
- }
- public String getName() {
- return name;
- }
- public void setName(String name) {
- this.name = name;
- }
- }
pom.xml:
- <dependencies>
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter</artifactId>
- </dependency>
-
- <dependency>
- <groupId>org.projectlombok</groupId>
- <artifactId>lombok</artifactId>
- <optional>true</optional>
- </dependency>
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-test</artifactId>
- <scope>test</scope>
- <exclusions>
- <exclusion>
- <groupId>org.junit.vintage</groupId>
- <artifactId>junit-vintage-engine</artifactId>
- </exclusion>
- </exclusions>
- </dependency>
- </dependencies>
RabbitReceiver.java(消费端,实际开发中可当成是service层,在controller层调用):
- package com.ue.conusmer;
-
- import com.rabbitmq.client.Channel;
- import com.ue.entity.Order;
- import org.springframework.amqp.rabbit.annotation.*;
- import org.springframework.amqp.support.AmqpHeaders;
- import org.springframework.messaging.Message;
- import org.springframework.messaging.handler.annotation.Headers;
- import org.springframework.messaging.handler.annotation.Payload;
- import org.springframework.stereotype.Component;
-
- import java.util.Map;
-
- @Component
- public class RabbitReceiver {
-
- @RabbitListener(bindings = @QueueBinding(
- value = @Queue(value = "queue-1",//指定消息队列名称
- durable="true"),
- exchange = @Exchange(value = "exchange-1",//指定交换机名称
- durable="true",
- type= "topic", //指定交换机的类型:主题交换机
- ignoreDeclarationExceptions = "true"),
- key = "springboot.*"//指定路由key
- )
- )
-
- @RabbitHandler
- public void onMessage(Message message, Channel channel) throws Exception {
- System.err.println("--------------------------------------");
- System.err.println("消费端Payload: " + message.getPayload());
- Long deliveryTag = (Long)message.getHeaders().get(AmqpHeaders.DELIVERY_TAG);
- //手工ACK
- channel.basicAck(deliveryTag, false);
- }
-
- /**
- * spring.rabbitmq.listener.order.queue.name=queue-2
- spring.rabbitmq.listener.order.queue.durable=true
- spring.rabbitmq.listener.order.exchange.name=exchange-1
- spring.rabbitmq.listener.order.exchange.durable=true
- spring.rabbitmq.listener.order.exchange.type=topic
- spring.rabbitmq.listener.order.exchange.ignoreDeclarationExceptions=true
- spring.rabbitmq.listener.order.key=springboot.*
- * @param order
- * @param channel
- * @param headers
- * @throws Exception
- */
- @RabbitListener(bindings = @QueueBinding(
- value = @Queue(value = "${spring.rabbitmq.listener.order.queue.name}",
- durable="${spring.rabbitmq.listener.order.queue.durable}"),
- exchange = @Exchange(value = "${spring.rabbitmq.listener.order.exchange.name}",
- durable="${spring.rabbitmq.listener.order.exchange.durable}",
- type= "${spring.rabbitmq.listener.order.exchange.type}",
- ignoreDeclarationExceptions = "${spring.rabbitmq.listener.order.exchange.ignoreDeclarationExceptions}"),
- key = "${spring.rabbitmq.listener.order.key}"
- )
- )
-
- @RabbitHandler
- public void onOrderMessage(@Payload Order order,
- Channel channel,
- @Headers Map<String, Object> headers) throws Exception {
- System.err.println("--------------------------------------");
- System.err.println("消费端order: " + order.getId());
- Long deliveryTag = (Long)headers.get(AmqpHeaders.DELIVERY_TAG);
- //手工ACK
- channel.basicAck(deliveryTag, false);
- }
- }
pom.xml:
- <dependencies>
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter</artifactId>
- </dependency>
-
- <dependency>
- <groupId>com.ue</groupId>
- <artifactId>rabbitmq-common</artifactId>
- <version>0.0.1-SNAPSHOT</version>
- </dependency>
-
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-test</artifactId>
- <scope>test</scope>
- <exclusions>
- <exclusion>
- <groupId>org.junit.vintage</groupId>
- <artifactId>junit-vintage-engine</artifactId>
- </exclusion>
- </exclusions>
- </dependency>
-
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-amqp</artifactId>
- </dependency>
- <dependency>
- <groupId>junit</groupId>
- <artifactId>junit</artifactId>
- <version>4.12</version>
- <scope>test</scope>
- </dependency>
- </dependencies>
application.properties:
- spring.rabbitmq.addresses=xxx:5672
- spring.rabbitmq.username=guest
- spring.rabbitmq.password=guest
- spring.rabbitmq.virtual-host=/
- spring.rabbitmq.connection-timeout=15000
-
- spring.rabbitmq.listener.simple.acknowledge-mode=manual
- spring.rabbitmq.listener.simple.concurrency=5
- spring.rabbitmq.listener.simple.max-concurrency=10
-
- spring.rabbitmq.listener.order.queue.name=queue-2
- spring.rabbitmq.listener.order.queue.durable=true
- spring.rabbitmq.listener.order.exchange.name=exchange-2
- spring.rabbitmq.listener.order.exchange.durable=true
- spring.rabbitmq.listener.order.exchange.type=topic
- spring.rabbitmq.listener.order.exchange.ignoreDeclarationExceptions=true
- spring.rabbitmq.listener.order.key=springboot.*
MainConfig.java:
- package com.ue;
-
- import org.springframework.context.annotation.ComponentScan;
- import org.springframework.context.annotation.Configuration;
-
- @Configuration
- @ComponentScan({"com.ue.*"})
- public class MainConfig {
-
- }
RabbitSender.java(生产端,实际开发中可当成是service层,在controller层调用):
- package com.ue.producer;
-
- import com.ue.entity.Order;
- import org.springframework.amqp.rabbit.connection.CorrelationData;
- import org.springframework.amqp.rabbit.core.RabbitTemplate;
- import org.springframework.amqp.rabbit.core.RabbitTemplate.ConfirmCallback;
- import org.springframework.amqp.rabbit.core.RabbitTemplate.ReturnCallback;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.messaging.Message;
- import org.springframework.messaging.MessageHeaders;
- import org.springframework.messaging.support.MessageBuilder;
- import org.springframework.stereotype.Component;
-
- import java.util.Map;
-
- @Component
- public class RabbitSender {
-
- //自动注入RabbitTemplate模板类
- @Autowired
- private RabbitTemplate rabbitTemplate;
-
- //回调函数: confirm确认,对于broker来说如果没有被签收或签收,都会给生产端发消息
- final ConfirmCallback confirmCallback = new ConfirmCallback() {
- @Override
- public void confirm(CorrelationData correlationData, boolean ack, String cause) {
- System.err.println("correlationData: " + correlationData);
- System.err.println("ack: " + ack);
- if(!ack){
- System.err.println("异常处理....");
- }
- }
- };
-
- //回调函数: return返回,路由键没有找到会进行下面的逻辑处理
- final ReturnCallback returnCallback = new ReturnCallback() {
- @Override
- public void returnedMessage(org.springframework.amqp.core.Message message, int replyCode, String replyText,
- String exchange, String routingKey) {
- System.err.println("return exchange: " + exchange + ", routingKey: "
- + routingKey + ", replyCode: " + replyCode + ", replyText: " + replyText);
- }
- };
-
- //发送消息方法调用: 构建Message消息
- public void send(Object message, Map<String, Object> properties) throws Exception {
- MessageHeaders mhs = new MessageHeaders(properties);
- Message msg = MessageBuilder.createMessage(message, mhs);
- rabbitTemplate.setConfirmCallback(confirmCallback);
- rabbitTemplate.setReturnCallback(returnCallback);
- //id + 时间戳 全局唯一
- CorrelationData correlationData = new CorrelationData("1234567890");
- rabbitTemplate.convertAndSend("exchange-1", "springboot.abc", msg, correlationData);
- }
-
- //发送消息方法调用: 构建自定义对象消息
- public void sendOrder(Order order) throws Exception {
- rabbitTemplate.setConfirmCallback(confirmCallback);
- rabbitTemplate.setReturnCallback(returnCallback);
- //id + 时间戳 全局唯一
- CorrelationData correlationData = new CorrelationData("0987654321");
- rabbitTemplate.convertAndSend("exchange-2", "springboot.def", order, correlationData);
- }
- }
pom.xml:
- <dependencies>
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter</artifactId>
- </dependency>
-
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-test</artifactId>
- <scope>test</scope>
- <exclusions>
- <exclusion>
- <groupId>org.junit.vintage</groupId>
- <artifactId>junit-vintage-engine</artifactId>
- </exclusion>
- </exclusions>
- </dependency>
-
- <dependency>
- <groupId>com.ue</groupId>
- <artifactId>rabbitmq-common</artifactId>
- <version>0.0.1-SNAPSHOT</version>
- </dependency>
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-amqp</artifactId>
- </dependency>
- <dependency>
- <groupId>junit</groupId>
- <artifactId>junit</artifactId>
- <version>4.12</version>
- <scope>test</scope>
- </dependency>
- </dependencies>
application.properties:
- spring.rabbitmq.addresses=xxx:5672
- spring.rabbitmq.username=guest
- spring.rabbitmq.password=guest
- spring.rabbitmq.virtual-host=/
- spring.rabbitmq.connection-timeout=15000
-
- spring.rabbitmq.publisher-confirms=true
- spring.rabbitmq.publisher-returns=true
- spring.rabbitmq.template.mandatory=true
MainConfig.java:
- package com.ue;
-
- import org.springframework.context.annotation.ComponentScan;
- import org.springframework.context.annotation.Configuration;
-
- @Configuration
- @ComponentScan({"com.ue.*"})
- public class MainConfig {
-
- }
RabbitmqSpringcloudProducerApplicationTests.java(测试类):
- package com.ue;
-
- import com.ue.entity.Order;
- import com.ue.producer.RabbitSender;
- import org.junit.Test;
- import org.junit.runner.RunWith;
- import org.springframework.beans.factory.annotation.Autowired;
- import org.springframework.boot.test.context.SpringBootTest;
- import org.springframework.test.context.junit4.SpringRunner;
-
- import java.text.SimpleDateFormat;
- import java.util.Date;
- import java.util.HashMap;
- import java.util.Map;
-
- @RunWith(SpringRunner.class)
- @SpringBootTest
- public class RabbitmqSpringcloudProducerApplicationTests {
-
- @Autowired
- private RabbitSender rabbitSender;
-
- private static SimpleDateFormat simpleDateFormat = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SSS");
-
- @Test
- public void testSender1() throws Exception {
- Map<String, Object> properties = new HashMap<>();
- properties.put("number", "12345");
- properties.put("send_time", simpleDateFormat.format(new Date()));
- rabbitSender.send("Hello RabbitMQ For Spring Boot!", properties);
- }
-
- @Test
- public void testSender2() throws Exception {
- Order order = new Order("001", "第一个订单");
- rabbitSender.sendOrder(order);
- }
- }
测试时先启动消费者,再运行测试的方法:
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。