blink
Java后端
·2023-11-11
在学习尚硅谷RabbitMQ课程的时候遇到了一些问题,是关于如何处理异步未确认消息的,网课上解决这个问题的方法是使用ConcurrentSkipListMap存储所有的消息,再在处理确认成功消息的时候将确认成功的消息从ConcurrentSkipListMap中删除,这样剩下的就是未确认的消息,但是他在处理批量消息时我有一些疑问, 他是靠headMap来获取键值严格小于deliveryTag的数据然后将其删除来保证剩下的都是未确认的消息 首先headMap获取的是键值严格小于deliveryTag的数据,他要删除确认数据的话应该是outstadingConfirms.headMap(deliveryTag 1)而不是 outstadingConfirms.headMap(deliveryTag) 然后他这样删除不会把未确认的消息也删掉吗? 如下例子 1, a -> 单个消息 确认删除 2 ,b -> 未确认数据 3 ,c -> 批量数据g 4 ,d -> 批量数据g ConcurrentNavigableMap<Long, String> confirmed = outstadingConfirms.headMap(deliveryTag); 通过这个方法应该获得的confirmed中应该存在 2,b : 3,c : 4,d这三组数据,因为outstadingConfirms中存储的是所有的消息,当然包括未确认的消息 那么使用confirmed.clear()删除的时候不会把未确认的消息也删除了吗? 希望大佬们帮忙解答一下 代码如下 //异步发布确认 public static void publishMessageAsync() throws Exception { Channel channel = RabbitMqUtils.getChannel(); //队列的声明 String queuename = UUID.randomUUID().toString(); channel.queueDeclare(queuename, true, false, false, null); //开启发布确认 channel.confirmSelect(); /** * 线程安全有序的一个哈希表 适用于高并发的情况下 * 1.轻松的将序号与消息进行关联 * * 2.轻松的批量删除条目 只要给到序号 * * 3.支持高并发(多线程) */ ConcurrentSkipListMap<Long, String> outstadingConfirms = new ConcurrentSkipListMap<>(); //消息确认成功 回调方法 ConfirmCallback ackCallback = (deliveryTag, multiple) -> { if(multiple) { //2:删除掉已经确认的消息 剩下的就是未确认的消息 //headMap是保留key小于deliverytag的部分tailMap是保留key大于deliverytag的部分 ConcurrentNavigableMap<Long, String> confirmed = outstadingConfirms.headMap(deliveryTag); //删除的confirmed集合内的数据,会让outstandingConfirms集合中的相应数据被删除 //headMep方法返回的是一个视图,对视图的操作会影响原来的 map,所以清空headMap没有问题 confirmed.clear(); } else { outstadingConfirms.remove(deliveryTag); } System.out.println("确认的消息:" deliveryTag); }; //消息确认失败 回调方法 /** * 1.消息的标记 * 2.是否为批量确认 */ ConfirmCallback nackCallback = (deliverTag, multiple) -> { //3:打印一下未确认的消息都有哪些 String message = outstadingConfirms.get(deliverTag); System.out.println("未确认的消息是" message "::::未确认的消息:" deliverTag); }; //准备消息的监听器,监听哪些消息成功了,哪些消息失败了 /** * 1.监听哪些消息成功了 * 2.监听哪些消息失败了 */ channel.addConfirmListener(ackCallback, nackCallback);//异步通知 //开始时间 long begin = System.currentTimeMillis(); //批量发送消息 for (int i = 0; i < MESSAGE_COUNT; i ) { String message = "消息" i; channel.basicPublish("", queuename, null, message.getBytes()); //1:此处记录下所有要发送的消息 消息的总和 outstadingConfirms.put(channel.getNextPublishSeqNo(), message); } //结束时间 long end = System.currentTimeMillis(); System.out.println("发布" MESSAGE_COUNT "个异步发布确认消息,耗时" (end - begin) "ms"); }
0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP