本节要点
- 理解消息丢失的几种情况
- 消费者如何保证幂等性
MQ消息的丢失
消息丢失的三种情况?
- 生产端丢失
- Broker丢失
- 消费端丢失

一、生产端

丢失的情况:
1、网络连接不通,业务成功了,但是mq发送失败
2、生产端业务报错,没有发送
解决方案:
1、同步发送,业务和发送消息在同一个事务中,抛出异常,消息确认失败,直接回滚
2、异步发送,业务和消息发送是异步的,业务成功;消息发送失败后重试,并且后继处理(补偿,记录消息手动处理)
同步做回滚;异步做补偿
二、MQ端

丢失情况:
1、刷盘的时候
异步刷盘:MQ收到消息还在内存里时候,就给生产端返回了确认。这是MQ还没有持久化,就会导致丢失。
解决方式:(SYNC_FLUSH)
2、集群复制时候
集群部署的时候,从库(Slave)还没有收到,主库(Master)的数据的时候,Master就宕机了。导致Slave接管的时候,没有刚才的那一条Master的数据
同步刷盘或主从同步复制效率相对较低,可根据实际业务权衡。
三、消费端
1、丢失场景:
自动回复ACK,业务处理失败,就会丢失

解决方式:
1、业务失败时直接抛出异常 或者 业务失败记录下来手动补充补偿
2、手动确认,保证同步处理,业务处理完毕以后再返回ack

存在的问题:
业务处理的时候,超过了MQ等待ACK的时间,MQ就认为没有消费成功,就会再次投递;但是业务最后是处理成功了!
幂等性问题
同一个业务,反复处理多次,但是只有处理一次的效果。(n)o

如何手动返回?
用原生的api类,代码如下:
package com.dw.listener;
import com.alibaba.fastjson.JSON;
import com.dw.message.TestMessage;
import jakarta.annotation.PostConstruct;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.springframework.stereotype.Component;
import java.nio.charset.StandardCharsets;
import java.util.List;
@Component
public class ManualAckConsumer {
@PostConstruct
public void init() throws MQClientException {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("dw-consumer-group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("topic-test", "*");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(
List<MessageExt> msgs,
ConsumeConcurrentlyContext context) {
for (MessageExt message : msgs) {
try {
String body = new String(message.getBody(), StandardCharsets.UTF_8);
TestMessage testMessage = JSON.parseObject(body, TestMessage.class);
System.out.println("接收到消息: " + testMessage);
//TO DO bussiness logic -> testMessage
// boolean flag = process(testMessage);
boolean flag = false;
if (flag) {
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 手动ACK
} else {
return ConsumeConcurrentlyStatus.RECONSUME_LATER; // 重试
}
} catch (Exception e) {
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER; // 重试
}
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
}
}
一直返回失败会发送16次,时间逐步递增
1. 递增延迟
- 从10秒开始,逐渐增加到2小时
- 这种设计是为了在系统故障时避免对系统造成过大压力
2. 第16次重试后
- 如果16次重试都失败,消息会进入死信队列(DLQ)
- 死信队列命名格式:
%DLQ%消费者组名 - 需要人工干预处理死信消息
如果要配置只需要在 broker.conf里增加
分别代表16个等级
messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
幂等性
1、什么是幂等性?
同一条数据,无论重复处理多少次都是一个结果
2、出现非幂等性的场景
1、手动ack还没有恢复的时候,网络就断开

2、网络没有断,但是处理的业务超时了

3、如何做到幂等性?
记录+判断
1、在redis中记录当前消息,msgid+状态(处理中,处理失败,处理成功)+超时时间

处理逻辑
- 消费者第一次收到消息时候,先把消息id和状态(处理中)存在redis里面
- 当业务处理完毕以后,把这一条消息的状态改为(处理成功)
- 就算是ack没有返回,导致mq重复投递消息,消费者拿到同样的消息的时候,在redis里判断这条消息id是否存在,如果存在就判断状态。如果是处理成功就直接返回ack=true,避免重复消费;如果是处理中,就不再处理业务直接返回ack=false,让MQ再次投递,就算多次循环,也不会导致重复消费。
实现这一目标的核心思路是:在消息处理前,先检查这个业务操作是否已经完成。如果已完成,则直接忽略本次重复消息;否则才执行真正的业务逻辑。
参考逻辑(可选)
package com.test.consumer;
import com.test.service.OrderInfoService;
import jakarta.annotation.PostConstruct;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.core.annotation.Order;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Component;
import java.nio.charset.StandardCharsets;
import java.util.List;
import static org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
@Component
public class ManualAck {
@Autowired
private RedisTemplate redisTemplate;
@Autowired
OrderInfoService orderInfoService;
@PostConstruct
public void init() throws MQClientException {//throws MQClientException {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("dw-consumer-group");
consumer.setNamesrvAddr("47.110.77.128:9876");
consumer.subscribe("topic-test", "*");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
for(MessageExt messageExt:list) {
try {
System.out.println("收到ack收到的消息------" + messageExt);
String msgId = messageExt.getMsgId();
String body = new String(messageExt.getBody(), StandardCharsets.UTF_8);
if(redisTemplate.opsForValue().get(msgId)==null){
System.out.println("---这是一个新的msgId,放入redis,状态处理中:"+msgId);
redisTemplate.opsForValue().set(msgId,"处理中");
orderInfoService.processBusiness(msgId); //业务处理
System.out.println("------------ack返回,未确认");
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}else{
String state = (String)redisTemplate.opsForValue().get(msgId);
System.out.println("---msgId已在redis中存在="+msgId+",state="+state);
if("处理中".equals(state)){
System.out.println("---msgId正在处理中,请稍后再试:"+msgId);
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}else if("处理完成".equals(state)){
System.out.println("---msgId处理完成,从redis中删除:"+msgId);
redisTemplate.delete(msgId);
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
}
}catch(Exception e){
//如果有异常就返回重发ACK
e.printStackTrace();
return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
}
//否者就告诉broker已经消费成功
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
}
}
25-3阶内容-2.2.1.4-MQ消息不丢失与幂等性
https://xiaochenblog.icu/archives/25-3jie-nei-rong-2.2.1.4-mqxiao-xi-bu-diu-shi-yu-mi-deng-xing
评论