栏目分类:
子分类:
返回
名师互学网用户登录
快速导航关闭
当前搜索
当前分类
子分类
实用工具
热门搜索
名师互学网 > IT > 前沿技术 > 大数据 > 大数据系统

05RabbitMq--如何保证MQ分布式事务的数据一致性,可靠性

05RabbitMq--如何保证MQ分布式事务的数据一致性,可靠性

生产者方面:

问题:生产者方面会出现消息投递不成功
解决:开启消息确认机制

spring.rabbitmq.publisher-/confirm/i-type=correlated

关于三种选值:

    none:默认值,不开启/confirm/icallback机制。correlated:开启/confirm/icallback,发布消息时,可以指定一个CorrelationData,会被保存到消息头中,消息投递到Broekr时触发生产者指定的/confirm/iCallback,这个值也会被返回,以进行对照处理,CorrelationData可以包含比较丰富的元信息进行回调逻辑的处理。无特殊需求,就设定为这个值。其一效果和correlated值一样能触发回调方法,其二用于发布消息成功后使用rabbitTemplate调用waitFor/confirm/is或waitFor/confirm/isOrDie方法等待broker节点返回发送结果,需求根据返回结果来判定下一步的逻辑,执行更复杂的业务。要注意的点是waitFor/confirm/isOrDie方法如果返回false则会关闭channel,则接下来无法发送消息。

配置类中配置:

@Configuration
public class RabbitConfig {
    @Autowired
    private RabbitTemplate rabbitTemplate;
    
    @PostConstruct
    public void setCallback() {
        
        rabbitTemplate.set/confirm/iCallback(new RabbitTemplate./confirm/iCallback() {
            
            @Override
            public void /confirm/i(CorrelationData correlationData, boolean ack, String cause) {
                if (ack) {
                    log.info("消息投递到交换机成功:[correlationData={}]",correlationData);
                } else {
                    log.error("消息投递到交换机失败:[correlationData={},原因:{}]", correlationData, cause);
                }
            }
        });
    }
}
消费者方面

问题:消费者方面会出现消费消息时失败的情况

解决:

    定义重试次数,失败时重新尝试消费消息手动获取ack,加上try/catchtry/catch+手动获取ack+死信队列
1.定义重试次数,失败时重新尝试消费消息

配置:

spring:
  # 项目名称
  application:
    name: rabbitmq-consumer
  # RabbitMQ服务配置
  rabbitmq:
    host: 127.0.0.1
    port: 5672
    username: guest
    password: guest
    listener:
      simple:
        # 重试机制
        retry:
          enabled: true #是否开启消费者重试
          max-attempts: 5 #最大重试次数
          initial-interval: 5000ms #重试间隔时间(单位毫秒)
          max-interval: 1200000ms #重试最大时间间隔(单位毫秒)
          multiplier: 2 #间隔时间乘子,间隔时间*乘子=下一次的间隔时间,最大不能超过设置的最大间隔时间
2.手动获取ack,加上try/catch

消费者模块中添加配置:

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: manual #手动应答
@Service
public class OrderConsumer {
    private DispatcherService dispatcherService;
    private int count = 1;
    @RabbitListener(queues = "save-order-queue")
    public void messageConsumer(String orderMsg, Channel channel, CorrelationData correlationData, @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
        try {
            System.out.println("收到mq的消息是:"+orderMsg+",count = "+count++);
            ;
            Order order = JSONObject.parseObject(orderMsg,Order.class);
            int orderId = order.getOrderId();
            boolean dispatcher = dispatcherService.dispatcher(orderId);
            if(dispatcher){
                System.out.println("消费者:ok");
            }else {
                System.out.println("消费者error");
            }
            channel.basicAck(tag,false);
        } catch (Exception e) {
            
            channel.basicNack(tag,false,false);
        }
    }
}
3.try/catch+手动获取ack+死信队列

基于解决办法2上,添加死信队列:

//死信队列管理
@Configuration
public class OrderProducer {
    @Bean
    public Queue deadQueue(){
        return new Queue("dead-queue",true);
    }
    @Bean
    public DirectExchange deadExchange() {
        return new DirectExchange("dead-exchange", true, false);
    }
    @Bean
    public Binding deadBindingExchange() {
        return BindingBuilder.bind(deadQueue()).to(deadExchange()).with("");
    }
    @Bean
    public Queue orderQueue(){
        Map args = new HashMap<>();
        args.put("x-dead-letter-exchange","dead-exchange");
        return new Queue("save-order-queue",true,false,false,args);
    }
    @Bean
    public DirectExchange orderExchange() {
        return new DirectExchange("save-order-exchange", true, false);
    }
    @Bean
    public Binding orderBingExchange() {
        return BindingBuilder.bind(orderQueue()).to(orderExchange()).with("");
    }
}
转载请注明:文章转载自 www.mshxw.com
本文地址:https://www.mshxw.com/it/735330.html
我们一直用心在做
关于我们 文章归档 网站地图 联系我们

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

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