本节要点
- 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 "发送成功";
}
}
21-3阶内容-2.2.1.1-MQ-SpringBoot集成(编写生产者)
https://xiaochenblog.icu/archives/21-3jie-nei-rong-2.2.1.1-mq-springbootji-cheng-bian-xie-sheng-chan-zhe
评论