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");
}
13
0
分享
操作
评论
问答助学
相关内容
0个评论
全部评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
