栏目分类:
子分类:
返回
名师互学网用户登录
快速导航关闭
当前搜索
当前分类
子分类
实用工具
热门搜索
名师互学网 > IT > 前沿技术 > 云计算 > 云平台

springboot 整合 RabbitMq 五种工作模式 消息手动确认和发送确认

云平台 更新时间: 发布时间: IT归档 最新发布 模块sitemap 名妆网 法律咨询 聚返吧 英语巴士网 伯小乐 网商动力

springboot 整合 RabbitMq 五种工作模式 消息手动确认和发送确认


整合maven:

       
            org.springframework.boot
            spring-boot-starter-amqp
            2.5.5
        

yml 配置:

rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    #开启发送到交换机确认callback
    #NONE值是禁用发布确认模式,是默认值
    #CORRELATED值是发布消息成功到交换器后会触发回调方法,如1示例
    #SIMPLE值经测试有两种效果,
    #其一效果和CORRELATED值一样会触发回调方法,#
    #其二在发布消息成功后使用rabbitTemplate调用waitForConfirms或waitForConfirmsOrDie方法等待broker节点返回发送结果,根据返回结果来判定下一步的逻辑,要注意的点是waitForConfirmsOrDie方法如果返回false则会关闭channel,则接下来无法发送消息到broker;
    publisher-confirm-type: correlated
    #开启发送到队列失败returnCallback
    publisher-returns: true
    #消息是否强制回退 如果此值为空才取publisher-returns值
    template:
      mandatory: true
    #开启手动确认
    listener:
      direct:
        acknowledge-mode: manual
      simple:
        acknowledge-mode: manual

 初始化Exchange 和Queue:

@Component
public class RabbitIniter {



    @Resource
    private AmqpAdmin amqpAdmin;


    @Resource
    private RabbitTemplate rabbitTemplate;

    @Resource
    private RabbitMqConfirmService rabbitMqConfirmService;

    @Bean
    public void initRabbitConfig(){
        rabbitTemplate.setConfirmCallback(rabbitMqConfirmService);
        rabbitTemplate.setReturnsCallback(rabbitMqConfirmService);
    }


    
    @Bean
    public void common() {
        amqpAdmin.declareQueue( new Queue("station_common"));
    }



    
    @Bean
    public void fanOut() {
        Queue stationFanOut = new Queue("station_fanOut");
        FanoutExchange fanoutExchange = new FanoutExchange("exchanges.fanout");
        amqpAdmin.declareQueue(stationFanOut);
        amqpAdmin.declareExchange(fanoutExchange);
        amqpAdmin.declareBinding(new Binding(
                stationFanOut.getName(),
                Binding.DestinationType.QUEUE,
                fanoutExchange.getName(),
                "",
                new HashMap<>(2)
                ));
    }


    
    @Bean
    public void route() {
        Queue stationRoute = new Queue("station_Route");
        DirectExchange directExchange = new DirectExchange("exchanges.route");
        amqpAdmin.declareQueue(stationRoute);
        amqpAdmin.declareExchange(directExchange);
        amqpAdmin.declareBinding(new Binding(
                stationRoute.getName(),
                Binding.DestinationType.QUEUE,
                directExchange.getName(),
                "route_exchange",
                new HashMap<>(4)
        ));
    }


    
    @Bean
    public void topic() {
        Queue stationTopic = new Queue("station_topic");
        TopicExchange topicExchange = new TopicExchange("exchanges.topic");
        amqpAdmin.declareQueue(stationTopic);
        amqpAdmin.declareExchange(topicExchange);
        amqpAdmin.declareBinding(new Binding(
                stationTopic.getName(),
                Binding.DestinationType.QUEUE,
                topicExchange.getName(),
                "*.log",
                new HashMap<>(6)
        ));
    }

 定义确认类实现


@Slf4j
@Component
public class RabbitMqConfirmService implements RabbitTemplate.ConfirmCallback, RabbitTemplate.ReturnsCallback {


    
    @Override
    public void confirm(CorrelationData correlationData, boolean isAck, String s) {
        if (isAck) {
            log.info("消息已发送到交换器 cause:{} - {}" , s , correlationData.toString());
        } else {
            log.info("消息未发送到交换器 cause:{} - {}" , s , correlationData.toString());
        }
    }
    
    @Override
    public void returnedMessage(ReturnedMessage returnedMessage) {
        log.info("消息被退回 {}" , returnedMessage.toString());
    }
}

 编写控制器:

@RestController
@RequestMapping("rabbit")
public class RabbitMqController {



    @Resource
    private RabbitTemplate rabbitTemplate;


    
    @GetMapping("common")
    public void sendMessage(){
        rabbitTemplate.convertAndSend("station_test","{name:123}");
    }



    
    @GetMapping("pub")
    public void pubsendMessage(){
        //消息确认时候 new CorrelationData()
        rabbitTemplate.convertAndSend("amq.fanout","station_test","{name:123}",new CorrelationData());
    }



    
    @GetMapping("direct")
    public void directsendMessage(){
        rabbitTemplate.convertAndSend("amq.direct","direct_route","{name:123}", new CorrelationData());
    }


    
    @GetMapping("topic")
    public void topicsendMessage(){
        rabbitTemplate.convertAndSend("amq_topic","error.log","{name:123}");
    }



}

 消息接收和确认

@Component
public class RabbtMqMessageReceiver {


    
    @RabbitListener(queues = "station_test")
    public void listenSimpleQueueMessage(String msg, Message message, Channel channel) throws InterruptedException, IOException {
        System.out.println(message);
        System.out.println("spring 消费者1接收到消息:【" + msg + "】");
        channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
    }


    



    @RabbitListener(queues = "station_test")
    public void listenSimpleQueueMessage2(String msg,Message message, Channel channel) throws InterruptedException, IOException {
        System.out.println("spring 消费者2接收到消息:【" + msg + "】");
        channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
    }




    @RabbitListener(queues = "station_common")
    public void listenSimpleQueueMessage3(String msg) throws InterruptedException {
        System.out.println("spring 消费者2接收到消息:【" + msg + "】");
    }



}

开启事务:

  @Bean
    public void initRabbitConfig(){
        rabbitTemplate.setConfirmCallback(rabbitMqConfirmService);
        rabbitTemplate.setReturnsCallback(rabbitMqConfirmService);
        rabbitTemplate.setChannelTransacted(true);
    }
   @Bean("rabbitTransactionManager")
    public RabbitTransactionManager rabbitTransactionManager(CachingConnectionFactory cachingConnectionFactory){
        return new RabbitTransactionManager(cachingConnectionFactory);
    }

开启事务就不能开启消息确认

    #开启事务的时候这个就不能设置了
    publisher-confirm-type: correlated
    
    @GetMapping("direct")
    @Transactional(rollbackFor = Exception.class,transactionManager = "rabbitTransactionManager")
    public void directsendMessage(){
        rabbitTemplate.convertAndSend("amq.direct","direct_route","{name:123}", new CorrelationData());
        rabbitTemplate.convertAndSend("amq.direct","direct_route","{name:456}", new CorrelationData());
    }

转载请注明:文章转载自 www.mshxw.com
本文地址:https://www.mshxw.com/it/896340.html
我们一直用心在做
关于我们 文章归档 网站地图 联系我们

版权所有 (c)2021-2022 MSHXW.COM

ICP备案号:晋ICP备2021003244-6号