复习回顾
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、消费端负责接收这个通知,并且调用阿里云的短信发送通知到用户
22-3阶内容-2.2.1.1-MQ(springboot)
https://xiaochenblog.icu/archives/22-3jie-nei-rong-2.2.1.1-mq-springboot
评论