复习回顾

1、什么是MQ

消息队列

2、核心作用(特点)

消峰、解耦、异步

3、 MQ的组成

NameServer、Broker、Producer、consomer

4、生成者

同步发送

异步发送

单向


延迟消息

顺序(生产者的发送方式)

事务(可选)

5、消费者

推和拉

6、同步和异步刷盘

保证数据不丢失,在接收到生产者消息的时候,同步或者异步刷盘

主从同步复制或异步复制

7、MQ概念

Topic,Message queue、offside、commin log

8、启动的流程

本节要点

  • SpringBoot集成MQ
  • 生产者和消费者
  • 自动ACK
  • 手动ACK(可选)

1、RocketMQ消息发送接收(重点)

生产者发送类型三种

  • 同步发送:同步发送会阻塞当前线程,直到收到 RocketMQ Broker 的确认(或超时)。这能确保你知道消息是否成功送达,但会牺牲一些性能。

  • **异步发送:**默认以异步方式发送消息(当 producer.sync=false 时)。这意味着 send 方法调用后会立即返回,发送结果由 RocketMQ Binder 在后台处理。

  • **单向发送:**只负责发送,不管消息是否发送成功。

    消费者消费消息分两种

  • **拉模式:**消费者主动去 Broker 上拉取消息。

  • **推模式:**消费者等待 Broker 把消息推送过来

使用方式:

  • 命令行
  • api调用
  • spring boot
  • spring cloud alibaba

2、RocketMQ集成SpringBoot

1、生产者

1)加依赖

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>2.3.0</version>
</dependency>

2)加配置

rocketmq:
  name-server: 127.0.0.1:9876  # RocketMQ命名服务器地址
  producer:
    group: my-producer-group   # 生产者组名
    send-message-timeout: 3000 # 发送消息超时时间,毫秒

3)编写发送对象(可选)

也可以发送一个字符串,但一般会用对象来发送,直接就是json到mq里

package com.dw.message;

import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;

import java.io.Serializable;

@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class TestMessage implements Serializable {
    private String msgId;
    private String msgContent;
}

4)编写生产者bean(service)

实际上就是用RocketMQTemplate发送

package com.dw.producer;

import com.dw.model.message.TestMessage;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

@Component
public class MqProducer {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    /**
     * 同步发送消息
     */
    public SendResult sendSyncMessage(String topic, String tag, TestMessage message) {
        // 同步发送消息
        String destination = topic + ":" + tag;
        return rocketMQTemplate.syncSend(destination, MessageBuilder.withPayload(message).build());
    }


    /**
     * 同步发送延时消息
     */
    public SendResult sendSyncDelayMessage(String topic, String message, int delayLevel) {
        return rocketMQTemplate.syncSend(topic,
                MessageBuilder.withPayload(message).build(),
                rocketMQTemplate.getProducer().getSendMsgTimeout(),
                delayLevel);
    }

    /**
     * 异步发送消息
     */
    public void sendAsyncMessage(String topic,String tag, TestMessage message) {
        String destination = topic + ":" + tag;
        rocketMQTemplate.asyncSend(destination, MessageBuilder.withPayload(message).build(), new SendCallback() {
      
            @Override
            public void onSuccess(SendResult sendResult) {
                System.out.println("异步消息发送成功: " + sendResult);
            }

            @Override
            public void onException(Throwable throwable) {
                System.out.println("异步消息发送失败: " + throwable.getMessage());
                // 处理发送失败逻辑
            }
        });
    }

    /**
     * 发送单向消息(不关心发送结果)
     */
    public void sendOnewayMessage(String topic, String message) {
        rocketMQTemplate.sendOneWay(topic, MessageBuilder.withPayload(message).build());
    }
  
  
    /**
     * 顺序发送
     */
    public void sendSyncOrderlyMessage(String topic, TestMessage message) {
        rocketMQTemplate.asyncSendOrderly(topic,
                MessageBuilder.withPayload(message).build(),
                "orderly",  // 这里设置hashKey - 使用业务ID如订单号 接收者也需要设置一样
                new SendCallback() {
                    @Override
                    public void onSuccess(SendResult sendResult) {
                        System.out.println("sendSyncOrderlyMessage-顺序消息-发送成功---");
                    }

                    @Override
                    public void onException(Throwable e) {
                        e.printStackTrace();
                    }
                });
    }
}
延迟级别 延迟时间 延迟级别 延迟时间
1 1秒 10 6分钟
2 5秒 11 7分钟
3 10秒 12 8分钟
4 30秒 13 9分钟
5 1分钟 14 10分钟
6 2分钟 15 20分钟
7 3分钟 16 30分钟
8 4分钟 17 1小时
9 5分钟 18 2小时

broker.conf 默认

# 默认配置,共18个等级
messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h

# 自定义配置示例:增加一个2天的等级,总共变成19个等级
messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h 2d

5)编写controller

package com.dw.controller;

import com.dw.model.message.TestMessage;
import com.dw.producer.MqProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

@RestController
@RequestMapping("/product")
public class MqTestController {
    @Autowired
    MqProducer mqProducer;

    /**
     * 发送同步消息
     * @return
     */
    @RequestMapping("/send/sync")
    public String sendSyncMessage() {
        TestMessage message = TestMessage.builder().msgId("1").msgContent("hello world").build();
        String str = mqProducer.sendSyncMessage("topic-test","Tag-A", message).toString();
        System.out.println(str);
        return "发送成功---"+str;
    }

    /**
     * 发送异步消息
     * @return
     */
    @RequestMapping("/send/async")
    public String sendAsyncMessage() {
        TestMessage message = TestMessage.builder().msgId("2").msgContent("嘻嘻哈哈").build();
        mqProducer.sendAsyncMessage("topic-test","Tag-B",message);
        return "发送成功";
    }

    /**
     * 发送一次性消息
     * @return
     */
    @RequestMapping("/send/oneway")
    public String sendOnewayMessage() {
        return "发送成功";
    }
}

2、消费者

1)加依赖

        <dependency>
            <groupId>org.apache.rocketmq</groupId>
            <artifactId>rocketmq-spring-boot-starter</artifactId>
            <version>2.3.0</version>
        </dependency>

2) 加配置

rocketmq:
  # 命名服务器地址(生产者和消费者都需要)
  #name-server:  47.99.194.36:9876
  name-server: 182.92.108.193:9876

3) 并发集群消费者

package com.dw.listener;

import com.dw.model.message.TestMessage;
import org.apache.rocketmq.spring.annotation.ConsumeMode;
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;

@Component
@RocketMQMessageListener(
        consumerGroup = "dw-consumer-group",
        topic = "topic-test",
        selectorExpression = "*",  //TAG
        //ackMode = AckMode.MANUAL,  // 设置为手动确认模式
        consumeMode = ConsumeMode.CONCURRENTLY, //# 消费模式:CONCURRENTLY(并发)或ORDERLY(顺序)
        messageModel = MessageModel.CLUSTERING  //# 消息模式:CLUSTERING(集群)或BROADCASTING(广播)
)

public class PushConsumer implements RocketMQListener<TestMessage> {

    @Override
    public void onMessage(TestMessage message) {
        try{
        // 处理业务逻辑
            System.out.println("Consumer 接收到消息: " + message);
        }catch(Exception e){
            e.printStackTrace();
        }
    }
}

4) 顺序消费

说明:接收者用consumeMode = ConsumeMode.ORDERLY, // 顺序消费模式

发送方(生产者):发送得使用增加HashKey,HashKey可以是自己定义,也可一个是某一类业务固定值

public void sendOrderlyMsg(String topic, String msg) {
        template.syncSendOrderly(topic, msg, "hashkey");
 }
package com.dw.listener;

import org.apache.rocketmq.spring.annotation.ConsumeMode;
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;

@Component
@RocketMQMessageListener(
        consumerGroup = "dw-consumer-group-orderly",
        topic = "topic-orderly",
        selectorExpression = "*", // 订阅所有Tag
        consumeMode = ConsumeMode.ORDERLY,  // 顺序消费模式
        messageModel = MessageModel.CLUSTERING
)
public class OrderlyConsumer implements RocketMQListener<String> {

    @Override
    public void onMessage(String message) {
        try {
            // 1. 处理消息(确保业务逻辑幂等)
            System.out.println("顺序消费消息: " + message);

            // 2. 模拟业务处理
        

        } catch (Exception e) {
            // 3. 如果处理失败,建议记录日志并人工干预
            // 注意:顺序消费下不要立即重试(会阻塞后续消息)
            e.printStackTrace();
        }
    }
}

完整配置

package com.dw.listener;

import com.dw.model.message.TestMessage;
import org.apache.rocketmq.spring.annotation.ConsumeMode;
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
@Component
@RocketMQMessageListener(
        consumerGroup = "dw-consumer-group",
        topic = "topic-test",
        selectorExpression = "*",    //TAG
        consumeMode = ConsumeMode.CONCURRENTLY, //# 消费模式:CONCURRENTLY(并发)或ORDERLY(顺序)
        messageModel = MessageModel.CLUSTERING  //# 消息模式:CLUSTERING(集群)或BROADCASTING(广播)
)
public class PushConsumer implements RocketMQListener<TestMessage> {

    @Override
    public void onMessage(TestMessage message) {
        try{
        // 处理业务逻辑
            System.out.println("Consumer1 接收到消息: " + message);
        }catch(Exception e){
            e.printStackTrace();
        }
    }
}

备注,如果需要顺序消费,修改consumeMode = ConsumeMode.ORDERLY,

//# 消费模式:CONCURRENTLY(并发)或ORDERLY(顺序)

一旦设置了循序消费,当前消费者只会消费指定的orderly消息队列

5)异常与重新投递

确认机制(自动确认)

onMessage 方法正常执行完成(没有抛出异常),RocketMQ 会自动确认成功:

3、手动ack(目前不支持RocketMQMessageListener)

但是可以通过抛出异常的方式让消息重新发送

 //throw new RuntimeException(e);  通过抛出异常可以让MQ从新发送

如果出现了异常就会,重复再次投递给消费者;并且每一次的时间逐渐加长,一共投递16次,就会进入死心队列(DLQ)

%DLQ%TOPIC

@Component
@RocketMQMessageListener(
        consumerGroup = "dw-consumer-group",   //消费者组
        topic = "topic-test",   //消费者监听的topic名称
        selectorExpression = "*",    //监听 topic-test 下的所有tag ,tag就是topic下面的分类
        consumeMode = ConsumeMode.CONCURRENTLY,   //并发消费和顺序消费
        messageModel = MessageModel.CLUSTERING,//集群消费,在有多个消费者消费同一个topic时,只有一个消费者拿到数据,不会导致消息重复处理,如果是BROADCASTING广播模式,就会导致消息被多个消费者一起消费
        //ackMode = AckMode.MANUAL  //没有手动模式
)
public class PushConsumer {
    public void onMessage(MessageExt message, Acknowledge acknowledge) {
        try {
            // 处理业务逻辑
            String messageBody = new String(message.getBody(), StandardCharsets.UTF_8);
            System.out.println("Consumer 接收到消息: " + messageBody);
          

            //acknowledge.ack();  没有手动模式,单可以用下面的 throw new RuntimeException(e),让消息重新发送
          
        } catch (Exception e) {
            e.printStackTrace();
            //throw new RuntimeException(e);  通过抛出异常可以让MQ从新发送
         	}
    }
}

4、手动ack用DefaultMQPushConsumer来实现(可选)

手动ack就用下面的对象实现

DefaultMQPushConsumer

参考代码

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("47.110.77.128: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);
                        //TO DO ...成功 flag = true;
                        boolean flag = true;
                        if (flag) {
                            System.out.println("确认成功:"+flag);
                            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 手动ACK
                        } else {
                            System.out.println("确认失败:"+flag);
                            return ConsumeConcurrentlyStatus.RECONSUME_LATER; // 重试
                        }
                    } catch (Exception e) {
                        e.printStackTrace();
                        return ConsumeConcurrentlyStatus.RECONSUME_LATER; // 重试
                    }
                }
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });

        consumer.start();
    }
}

流程演示:

1、创建一个订单

2、创建订单以后,发送一个短信通知的消息 (生产者)

3、消费端负责接收这个通知,并且调用阿里云的短信发送通知到用户