当前位置:   article > 正文

SpringBoot配置RabbitMq消息队列_springboot rabbitmq listener 配置

springboot rabbitmq listener 配置

RabbitMq配置依赖

compile ‘org.springframework.amqp:spring-rabbit:1.7.5.RELEASE’

RabbitMq生产配置

package com.autoyol.car.study.service.message;

public class RabbitMqConfirmProducer implements ConfirmCallback {
    private static final Logger log = LoggerFactory.getLogger(RabbitMqConfirmProducer.class);
    private RabbitTemplate confirmRabbitTemplate;

    public RabbitMqConfirmProducer(RabbitTemplate confirmRabbitTemplate) {
        this.confirmRabbitTemplate = confirmRabbitTemplate;
    }

    public void sendTopicMsg(String exchange, String routeKey, Object message) {
        String correlationDataId = UUID.randomUUID().toString();
        if (message == null) {
            return;
        }
        this.confirmRabbitTemplate.setConfirmCallback(this);
        CorrelationData correlationData = new CorrelationData(correlationDataId);
        this.confirmRabbitTemplate.convertAndSend(exchange, routeKey, message, correlationData);
    }

    @Override
    public void confirm(CorrelationData correlationData, boolean ack, String cause) {
        if (ack) {
            log.info("messageId is " + correlationData.getId() + ",send message success");
        } else {
            log.error("messageId is " + correlationData.getId() + ",send message failed:" + cause);
        }
    }

}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28

RabbitMq消费配置

package com.autoyol.car.study.service.message;

/**
 * @Description: 消息序列化
 * @menu: FastJsonMessageConverter
 * @Date: 2021/2/25
 **/
public class JsonMessageConverter extends AbstractMessageConverter {
    private static Logger log = LoggerFactory.getLogger(JsonMessageConverter.class);
    public static final String DEFAULT_CHARSET = "UTF-8";

    public JsonMessageConverter() {
    }

    @Override
    public Object fromMessage(Message message) throws MessageConversionException {
        String json = "";
        try {
            json = new String(message.getBody(), DEFAULT_CHARSET);
        } catch (UnsupportedEncodingException e) {
            e.printStackTrace();
        }
        return JsonUtil.fromJson(json, Object.class);
    }

    @Override
    protected Message createMessage(Object objectToConvert, MessageProperties messageProperties) throws MessageConversionException {
        byte[] bytes;
        try {
            String jsonString = JsonUtil.toJson(objectToConvert);
            bytes = jsonString.getBytes(DEFAULT_CHARSET);
        } catch (UnsupportedEncodingException e) {
            throw new MessageConversionException("Failed to convert Message content", e);
        }
        messageProperties.setContentType("application/json");
        messageProperties.setContentEncoding(DEFAULT_CHARSET);
        if (bytes != null) {
            messageProperties.setContentLength((long) bytes.length);
        }
        return new Message(bytes, messageProperties);
    }
}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28
  • 29
  • 30
  • 31
  • 32
  • 33
  • 34
  • 35
  • 36
  • 37
  • 38
  • 39
  • 40

package com.autoyol.car.study.service.message;

@Configuration
@Data
public class RabbitMqConsumeConfig {
    private final static Logger LOGGER = LoggerFactory.getLogger(RabbitMqConsumeConfig.class);
    private final static Integer prefetchCount = 10;
    @Bean
    public MessageConverter jsonMessageConverter() {
        return new JsonMessageConverter();
    }
    @Bean
    public RabbitTemplate rabbitTemplate(CachingConnectionFactory connectionFactory) {
        RabbitTemplate template = new RabbitTemplate(connectionFactory);
        return template;
    }
    @Bean
    @Qualifier(value = "confirmRabbitTemplate")
    @Scope("prototype")
    public RabbitTemplate confirmRabbitTemplate(CachingConnectionFactory connectionFactory) {
        connectionFactory.setPublisherConfirms(true);
        RabbitTemplate template = new RabbitTemplate(connectionFactory);
        template.setMessageConverter(this.jsonMessageConverter());
        return template;
    }
    @Bean
    public RabbitMqConfirmProducer rabbitMqConfirmProducer(CachingConnectionFactory connectionFactory) {
        RabbitMqConfirmProducer template = new RabbitMqConfirmProducer(confirmRabbitTemplate(connectionFactory));
        return template;
    }
    @Bean("rabbitListenerContainerFactory")
    public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(ConnectionFactory connectionFactory) {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(connectionFactory);
        // 并发消费者数量
        factory.setConcurrentConsumers(10);
        factory.setMaxConcurrentConsumers(50);
        factory.setAcknowledgeMode(AcknowledgeMode.AUTO);
        return factory;
    }
    @Bean("simpleRabbitListenerSimpleFactory")
    public SimpleRabbitListenerContainerFactory simpleRabbitListenerSimpleFactory(ConnectionFactory connectionFactory) {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(connectionFactory);
        factory.setMessageConverter(this.jsonMessageConverter());
        factory.setConcurrentConsumers(prefetchCount);
        // 并发消费者数量
        factory.setConcurrentConsumers(10);
        factory.setMaxConcurrentConsumers(50);
        factory.setAcknowledgeMode(AcknowledgeMode.AUTO);
        return factory;
    }
    @Bean("confirmRabbitListenerContainerFactory")
    public SimpleRabbitListenerContainerFactory confirmRabbitListenerContainerFactory(ConnectionFactory connectionFactory) {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(connectionFactory);
        factory.setMessageConverter(jsonMessageConverter());
        // 并发消费者数量
        factory.setConcurrentConsumers(10);
        factory.setMaxConcurrentConsumers(50);
        /** 设置当rabbitmq收到nack/reject确认信息时的处理方式,设为true,扔回queue头部,设为false,丢弃。*/
        factory.setDefaultRequeueRejected(true);
        factory.setAcknowledgeMode(AcknowledgeMode.MANUAL);
        factory.setAdviceChain(
                RetryInterceptorBuilder
                        .stateless()
                        .recoverer(new RejectAndDontRequeueRecoverer())
                        .retryOperations(rabbitRetryTemplate())
                        .build()
        );
        return factory;
    }
    @Bean
    public RetryTemplate rabbitRetryTemplate() {
        RetryTemplate retryTemplate = new RetryTemplate();
        // 设置监听(不是必须)
        retryTemplate.registerListener(new RetryListener() {
            @Override
            public <T, E extends Throwable> boolean open(RetryContext retryContext, RetryCallback<T, E> retryCallback) {
                // 执行之前调用 (返回false时会终止执行)
                return true;
            }
            @Override
            public <T, E extends Throwable> void close(RetryContext retryContext, RetryCallback<T, E> retryCallback, Throwable throwable) {
                // 重试结束的时候调用 (最后一次重试 )
            }
            @Override
            public <T, E extends Throwable> void onError(RetryContext retryContext, RetryCallback<T, E> retryCallback, Throwable throwable) {
                //  异常 都会调用
                LOGGER.info("当前队列消费重试次数:{},异常信息:", retryContext.getRetryCount(), throwable);
            }
        });
        // 个性化处理异常和重试 (不是必须)
        Map<Class<? extends Throwable>, Boolean> retryableExceptions = new HashMap<>();
        //设置重试异常和是否重试
        retryableExceptions.put(AmqpException.class, true);
        //设置重试次数和要重试的异常
        SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(5, retryableExceptions);
        retryTemplate.setBackOffPolicy(backOffPolicyByProperties());
        retryTemplate.setRetryPolicy(retryPolicy);
        retryTemplate.setThrowLastExceptionOnExhausted(true);
        return retryTemplate;
    }
    @Bean
    public ExponentialBackOffPolicy backOffPolicyByProperties() {
        ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy();
        // 重试间隔
        backOffPolicy.setInitialInterval(10 * 1000);
        // 重试最大间隔
        backOffPolicy.setMaxInterval(300 * 1000);
        // 重试间隔乘法策略
        backOffPolicy.setMultiplier(3);
        return backOffPolicy;
    }
}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
  • 23
  • 24
  • 25
  • 26
  • 27
  • 28
  • 29
  • 30
  • 31
  • 32
  • 33
  • 34
  • 35
  • 36
  • 37
  • 38
  • 39
  • 40
  • 41
  • 42
  • 43
  • 44
  • 45
  • 46
  • 47
  • 48
  • 49
  • 50
  • 51
  • 52
  • 53
  • 54
  • 55
  • 56
  • 57
  • 58
  • 59
  • 60
  • 61
  • 62
  • 63
  • 64
  • 65
  • 66
  • 67
  • 68
  • 69
  • 70
  • 71
  • 72
  • 73
  • 74
  • 75
  • 76
  • 77
  • 78
  • 79
  • 80
  • 81
  • 82
  • 83
  • 84
  • 85
  • 86
  • 87
  • 88
  • 89
  • 90
  • 91
  • 92
  • 93
  • 94
  • 95
  • 96
  • 97
  • 98
  • 99
  • 100
  • 101
  • 102
  • 103
  • 104
  • 105
  • 106
  • 107
  • 108
  • 109
  • 110
  • 111
  • 112
  • 113

RabbitMq使用样例

package com.autoyol.car.study.service.message;

@Component
public class TemplateRabbitConsumer {

    private final static Logger LOGGER = LoggerFactory.getLogger(TemplateRabbitConsumer.class);
    public final static String topic = "weixin_topic";
    public final static String topic_key = "weixin_key";

    @RabbitListener(bindings = @QueueBinding(
            value = @Queue(value = "weixin_queue", durable = "true", autoDelete = "true"),
            exchange = @Exchange(value = topic, type = ExchangeTypes.TOPIC),
            key = topic_key), containerFactory = "confirmRabbitListenerContainerFactory")
    public void onMessageOfOfRabbit(Message message, Channel channel) {
        try {
            LOGGER.info("onMessageOfOfRabbit->>当前消息数据:{}", new String(message.getBody()));
            User user = JsonUtil.fromJson(new String(message.getBody()), User.class);
            LOGGER.info("onMessageOfOfRabbit->>当前车辆编号:{}", JsonUtil.toJson(user));
            channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
        } catch (Exception e) {
            LOGGER.info("exception is ", e);
        }
    }
}
  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
  • 11
  • 12
  • 13
  • 14
  • 15
  • 16
  • 17
  • 18
  • 19
  • 20
  • 21
  • 22
声明:本文内容由网友自发贡献,不代表【wpsshop博客】立场,版权归原作者所有,本站不承担相应法律责任。如您发现有侵权的内容,请联系我们。转载请注明出处:https://www.wpsshop.cn/w/我家小花儿/article/detail/483413
推荐阅读
相关标签
  

闽ICP备14008679号