本节要点

  • 理解消息丢失的几种情况
  • 消费者如何保证幂等性

MQ消息的丢失

消息丢失的三种情况?

  • 生产端丢失
  • Broker丢失
  • 消费端丢失

image-20250814181036934

一、生产端

image-20250814183937863

丢失的情况:

1、网络连接不通,业务成功了,但是mq发送失败

2、生产端业务报错,没有发送

解决方案:

1、同步发送,业务和发送消息在同一个事务中,抛出异常,消息确认失败,直接回滚

2、异步发送,业务和消息发送是异步的,业务成功;消息发送失败后重试,并且后继处理(补偿,记录消息手动处理)

同步做回滚;异步做补偿

二、MQ端

image-20250814185352915

丢失情况:

1、刷盘的时候

异步刷盘:MQ收到消息还在内存里时候,就给生产端返回了确认。这是MQ还没有持久化,就会导致丢失。

解决方式:(SYNC_FLUSH)

2、集群复制时候

集群部署的时候,从库(Slave)还没有收到,主库(Master)的数据的时候,Master就宕机了。导致Slave接管的时候,没有刚才的那一条Master的数据

同步刷盘或主从同步复制效率相对较低,可根据实际业务权衡。

三、消费端

1、丢失场景:

自动回复ACK,业务处理失败,就会丢失

image-20251022195502816

解决方式:

1、业务失败时直接抛出异常 或者 业务失败记录下来手动补充补偿

2、手动确认,保证同步处理,业务处理完毕以后再返回ack

image-20251022195656161

存在的问题:

业务处理的时候,超过了MQ等待ACK的时间,MQ就认为没有消费成功,就会再次投递;但是业务最后是处理成功了!

幂等性问题

同一个业务,反复处理多次,但是只有处理一次的效果。(n)o

image-20250815094427387

如何手动返回?

用原生的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还没有恢复的时候,网络就断开

image-20251022200900237

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

image-20251022200841743

3、如何做到幂等性?

记录+判断

1、在redis中记录当前消息,msgid+状态(处理中,处理失败,处理成功)+超时时间

image-20250815094905705

处理逻辑

  • 消费者第一次收到消息时候,先把消息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();

    }


}