|
@@ -0,0 +1,124 @@
|
|
|
|
|
+package com.ydtech.config;
|
|
|
|
|
+
|
|
|
|
|
+import com.ydtech.constants.RabbitMqEnum;
|
|
|
|
|
+import org.slf4j.Logger;
|
|
|
|
|
+import org.slf4j.LoggerFactory;
|
|
|
|
|
+import org.springframework.amqp.core.*;
|
|
|
|
|
+import org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory;
|
|
|
|
|
+import org.springframework.amqp.rabbit.connection.ConnectionFactory;
|
|
|
|
|
+import org.springframework.amqp.rabbit.connection.CorrelationData;
|
|
|
|
|
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
|
|
|
|
|
+import org.springframework.boot.autoconfigure.amqp.SimpleRabbitListenerContainerFactoryConfigurer;
|
|
|
|
|
+import org.springframework.context.annotation.Bean;
|
|
|
|
|
+import org.springframework.context.annotation.Configuration;
|
|
|
|
|
+
|
|
|
|
|
+@Configuration
|
|
|
|
|
+public class RabbitMQConfig {
|
|
|
|
|
+
|
|
|
|
|
+ /**
|
|
|
|
|
+ * 消费者数量,默认10
|
|
|
|
|
+ */
|
|
|
|
|
+ public static final int DEFAULT_CONCURRENT = 10;
|
|
|
|
|
+
|
|
|
|
|
+ /**
|
|
|
|
|
+ * 每个消费者获取最大投递数量 默认50
|
|
|
|
|
+ */
|
|
|
|
|
+ public static final int DEFAULT_PREFETCH_COUNT = 50;
|
|
|
|
|
+
|
|
|
|
|
+ private final static Logger log = LoggerFactory.getLogger(RabbitMQConfig.class);
|
|
|
|
|
+
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ //设置队列持久化 第二个参数 保证数据不丢失
|
|
|
|
|
+ Queue queue() {
|
|
|
|
|
+ return new Queue(RabbitMqEnum.ORDER_QUEUE.getCode(), true);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ //设置队列持久化 第二个参数 保证数据不丢失
|
|
|
|
|
+ Queue queueFail() {
|
|
|
|
|
+ return new Queue(RabbitMqEnum.ORDER_QUEUE_FAIL.getCode(), true);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ //设置队列持久化 第二个参数 保证数据不丢失
|
|
|
|
|
+ Queue queueSign() {
|
|
|
|
|
+ return new Queue(RabbitMqEnum.ORDER_QUEUE_SIGN.getCode(), true);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ //设置队列持久化 第二个参数 保证数据不丢失
|
|
|
|
|
+ DirectExchange exchange() {
|
|
|
|
|
+ return new DirectExchange(RabbitMqEnum.ORDER_EXCHANGE.getCode(),true,false);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ //设置队列持久化 第二个参数 保证数据不丢失
|
|
|
|
|
+ DirectExchange exchangeFail() {
|
|
|
|
|
+ return new DirectExchange(RabbitMqEnum.ORDER_EXCHANGE_FAIL.getCode(),true,false);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ //设置队列持久化 第二个参数 保证数据不丢失
|
|
|
|
|
+ DirectExchange exchangeSign() {
|
|
|
|
|
+ return new DirectExchange(RabbitMqEnum.ORDER_EXCHANGE_SIGN.getCode(),true,false);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ Binding binding(Queue queue, DirectExchange exchange) {
|
|
|
|
|
+ return BindingBuilder.bind(queue).to(exchange).with(RabbitMqEnum.ORDER_ROUTINGKEY.getCode());
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ Binding bindingFail(Queue queue, DirectExchange exchange) {
|
|
|
|
|
+ return BindingBuilder.bind(queue).to(exchange).with(RabbitMqEnum.ORDER_ROUTINGKEY_FAIL.getCode());
|
|
|
|
|
+ }
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ Binding bindingSign(Queue queue, DirectExchange exchange) {
|
|
|
|
|
+ return BindingBuilder.bind(queue).to(exchange).with(RabbitMqEnum.ORDER_ROUTINGKEY_SIGN.getCode());
|
|
|
|
|
+ }
|
|
|
|
|
+ /**
|
|
|
|
|
+ * @version
|
|
|
|
|
+ * @author: hxl
|
|
|
|
|
+ * @Date: 2024/11/5 10:27
|
|
|
|
|
+ * @Description: 手动确认保证消息消费成功
|
|
|
|
|
+ */
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ public RabbitTemplate createRabbitTemplate(ConnectionFactory connectionFactory) {
|
|
|
|
|
+ RabbitTemplate rabbitTemplate = new RabbitTemplate();
|
|
|
|
|
+ rabbitTemplate.setConnectionFactory(connectionFactory);
|
|
|
|
|
+ rabbitTemplate.setMandatory(true);
|
|
|
|
|
+ // 设置配置回调函数
|
|
|
|
|
+ rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
|
|
|
|
|
+ @Override
|
|
|
|
|
+ public void confirm(CorrelationData correlationData, boolean ack, String cause) {
|
|
|
|
|
+ log.info("ConfirmCallback: " + "相关数据:" + correlationData);
|
|
|
|
|
+ log.info("ConfirmCallback: " + "确认情况:" + ack);
|
|
|
|
|
+ log.info("ConfirmCallback: " + "原因:" + cause);
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
|
|
+ rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
|
|
|
|
|
+ @Override
|
|
|
|
|
+ public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
|
|
|
|
|
+ log.info("ReturnCallback: " + "消息:" + message);
|
|
|
|
|
+ log.info("ReturnCallback: " + "回应码:" + replyCode);
|
|
|
|
|
+ log.info("ReturnCallback: " + "回应信息:" + replyText);
|
|
|
|
|
+ log.info("ReturnCallback: " + "交换机:" + exchange);
|
|
|
|
|
+ log.info("ReturnCallback: " + "路由键:" + routingKey);
|
|
|
|
|
+ }
|
|
|
|
|
+ });
|
|
|
|
|
+ return rabbitTemplate;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+
|
|
|
|
|
+ @Bean
|
|
|
|
|
+ public SimpleRabbitListenerContainerFactory pointTaskContainerFactory(SimpleRabbitListenerContainerFactoryConfigurer configurer, ConnectionFactory connectionFactory) {
|
|
|
|
|
+ SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
|
|
|
|
|
+ factory.setPrefetchCount(DEFAULT_PREFETCH_COUNT);
|
|
|
|
|
+ factory.setConcurrentConsumers(DEFAULT_CONCURRENT);
|
|
|
|
|
+ factory.setAcknowledgeMode(AcknowledgeMode.MANUAL);
|
|
|
|
|
+ configurer.configure(factory, connectionFactory);
|
|
|
|
|
+ return factory;
|
|
|
|
|
+ }
|
|
|
|
|
+}
|