springboot集成rocketmq基础
回顾一下springboot集成rocketmq的一些用法,实例进行测试。
目录
rocketmq简介
基础概念
producer:生产者,发送消息到broker
消息发送的三种方式:同步(阻塞等待)、异步(发送结果回调)、单向(不等待broker响应)
consumer:消费者,从broker拉取消息进行消费
消息拉取的两种方式:pull、push
broker:消息服务器,处理生产者、消费者的请求,接收、存储、转发消息,存储消息元数据
nameserver:rocketmq的注册中心,维护broker的路由信息,生产者和消费者通过nameserver获取broker的地址信息
topic:主题
tag:标签,消息的二级分类,消费时方便过滤消息
message queue:消息队列,每个topic由多个queue组成
message:消息,包含topic、tag、内容、唯一标识key等信息
group:分组,生产者和消费者均可按业务进行分组
消息的发送与消费
发送流程:
生产者生产消息,指定topic
通过nameserver获取topic对应的broker路由信息
根据路由信息发送消息到broker
broker接收到消息后,持久化到commitLog
消费流程:
消费者启动时向nameserver注册自己,订阅topic
通过nameserver获取topic对应的broker路由信息
从broker中获取消息进行消费,消费完成后手动或自动确认
消息的可靠性和高可用性
消息的可靠性:
生产者可靠:发送消息后,需等待broker确认
消息存储可靠:broker将消息持久化到commitLog,确保消息不丢失
消费者可靠:消费时需要消费者确认,消费失败将自动重试
消息存储机制:
commitLog:追加的方式快速写入、批量读取
consumerQueue:为每个topic创建一到多个queue,存储指向commitLog的偏移量
indexFile:各类索引文件,如消息Id、key,便于快速检索消息
broker主从架构:
master负责处理读写请求,slave仅备份数据
master故障时,slave升级为master
消息生产消费性能优化:
生产者批量发送消息
消费者多线程并发处理消息
消息压缩技术
docker安装rocketmq
官网:https://rocketmq.apache.org/zh/docs/
下拉镜像:配置文件路径指定了版本号,所以下拉指定版本
docker pul apache/rocketmq:5.3.2
创建共享网络,便于通信
docker network create rmq_network
启动nameserver
docker run -d -p 9876:9876 --name my_rmq_server --network rmq_network apache/rocketmq:5.3.2 sh mqnamesrv
代码中连接使用9876
rocketmq.name-server=localhost:9876
启动broker
本地某个目录下创建broker配置文件broker.conf:
namesrvAddr = host.docker.internal:9876
brokerClusterName = DefaultCluster
# 节点名称
brokerName = broker-a
# 0 means master other means slave
brokerId = 0
brokerIP1 = host.docker.internal
brokerRole = ASYNC_MASTER
flushDiskType = ASYNC_FLUSH
# delete commit log at 04:00
deleteWhen = 04
fileReservedTime = 72
autoCreateTopicEnable=true
autoCreateSubscriptionGroup=true
启动broker,指定配置文件映射(本地配置文件路径D:\\Cache\\Docker\\rocketmq\\broker.conf)
docker run -d -p 10912:10912 -p 10911:10911 -p 10909:10909 --name my_rmq_broker ^
--network rmq_network ^
-v D:\\Cache\\Docker\\rocketmq\\broker.conf:/home/rocketmq/rocketmq-5.3.2/conf/broker.conf ^
apache/rocketmq:5.3.2 sh mqbroker -c /home/rocketmq/rocketmq-5.3.2/conf/broker.conf
安装dashboard
docker run -d -p 8099:8080 --name my_rmq_dashboard --network rmq_network -e "JAVA_OPTS=-Drocketmq.namesrv.addr=host.docker.internal:9876" apacherocketmq/rocketmq-dashboard
访问dashboard,查看各种信息:
http://localhost:8099/

springboot集成rocketmq
原生rocketmq客户端方式
添加依赖
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.2.3</version>
</dependency>
application.properties配置
rocketmq.name-server=localhost:9876
rocketmq.producer.group=producer_group_demo
简单Demo
topic:
首先在dashboard创建普通topic:normal_topic
启动消费者
public static void main(String[] args) {
// pullConsumer is deprecated, use push instead
DefaultMQPushConsumer pushConsumer = new DefaultMQPushConsumer("consumer_group_demo");
pushConsumer.setNamesrvAddr("localhost:9876");
pushConsumer.subscribe("normal_topic", "*");
pushConsumer.registerMessageListener((MessageListenerConcurrently)
(list, context) -> {
log.info("receive message size: {}", list.size());
for (MessageExt messageExt : list) {
log.info("receive message: {}", messageExt);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
log.info("consumer started");
pushConsumer.start();
}
启动生产者
@SneakyThrows
public static void main(String[] args) {
DefaultMQProducer producer = new DefaultMQProducer("producer_group_demo");
producer.setNamesrvAddr("localhost:9876");
producer.start();
for (int i = 0; i < 10; i++) {
Message message = new Message();
message.setTopic("normal_topic");
message.setTags("tag");
message.setBody(("normal_topic_msg_" + i).getBytes());
SendResult sendResult = producer.send(message);
log.info("sendResult: {}", sendResult);
}
producer.shutdown();
}
普通消息
topic:
首先在dashboard创建普通topic:normal_topic
消费者A:使用泛型接收消息内容
@Component
@Slf4j
@RocketMQMessageListener(topic = "normal_topic", consumerGroup = "normal_consumer_a")
public class NormalConsumerA implements RocketMQListener<String> {
@Override
@SneakyThrows
public void onMessage(String message) {
log.info("NormalConsumerA onMessage: {}", message);
Thread.sleep(1000);
}
}
消费者B:指定tag=b,按业务逻辑进行分类处理
@Component
@Slf4j
@RocketMQMessageListener(topic = "normal_topic", consumerGroup = "normal_consumer_b",
selectorType = SelectorType.TAG, selectorExpression = "b")
public class NormalConsumerB implements RocketMQListener<String> {
@Override
@SneakyThrows
public void onMessage(String message) {
log.info("NormalConsumerB thread {} onMessage: {}",
Thread.currentThread().getName(), message);
Thread.sleep(1000);
}
}
消费者C:手动确认,需要使用MessageExt
@Component
@Slf4j
@RocketMQMessageListener(topic = "normal_topic", consumerGroup = "normal_consumer_ack")
public class NormalConsumerAck implements RocketMQReplyListener<MessageExt, String> {
@Override
@SneakyThrows
public String onMessage(MessageExt message) {
String msg = new String(message.getBody());
log.info("NormalConsumerAck onMessage: {}", msg);
Thread.sleep(1000);
// 返回非null标识成功 null或异常触发重试
return "success";
}
}
生产者:部分代码
@Resource
private RocketMQTemplate rocketMQTemplate;
String msg = "normal_topic_msg_" + LocalDateTime.now();
// 同步发送 - 指定tag a
SendResult sendResultA = rocketMQTemplate.syncSend(topic, msg);
log.info("sendResultA: {}", sendResultA);
// 同步发送 - 指定tag b
SendResult sendResultB = rocketMQTemplate.syncSend(topic + ":b",
MessageBuilder.withPayload(msg).build());
log.info("sendResultB: {}", sendResultB);
// 异步发送 回调
rocketMQTemplate.asyncSend(topic, MessageBuilder.withPayload(msg).build(),
new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("SendCallback sendResult: {}", sendResult);
}
@Override
public void onException(Throwable throwable) {
log.error("SendCallback throwable: {}", throwable.getMessage());
}
});
// 异步发送 - spring 简化
rocketMQTemplate.convertAndSend(topic, msg);
log.info("convertAndSend: {}", msg);
// 单向发送 - 只发送 不等待broker确认响应
rocketMQTemplate.sendOneWay(topic, msg);
顺序消息
topic:
首先在dashboard创建顺序topic:orderly_topic
消费者
@Component
@Slf4j
@RocketMQMessageListener(topic = "orderly_topic", consumerGroup = "orderly_consumer_group",
consumeMode = ConsumeMode.ORDERLY, consumeThreadNumber = 1)
public class OrderlyConsumer implements RocketMQListener<String> {
@Override
@SneakyThrows
public void onMessage(String message) {
int cost = new Random().nextInt(5);
log.info("OrderlyConsumer thread {} onMessage: {} cost: {}",
Thread.currentThread().getName(), message, cost);
TimeUnit.SECONDS.sleep(cost);
}
}
生产者
for (int i = 0; i < 10; i++) {
String msg = "orderly_topic_msg_" + i;
SendResult sendResult = rocketMQTemplate.syncSendOrderly(topic, msg, "userId");
log.info("sendResult: {}", sendResult);
}
延时/定时消息
topic:
首先在dashboard创建延时topic:delay_topic
消费者
@Component
@Slf4j
@RocketMQMessageListener(topic = "delay_topic", consumerGroup = "delay_consumer_group")
public class DelayConsumer implements RocketMQListener<DelayMsg> {
@Override
public void onMessage(DelayMsg message) {
long delay = (System.currentTimeMillis() - message.getTime()) / 1000;
log.info("DelayConsumer delayLevel: {} delay: {}", message.getDelayLevel(), delay);
}
}
生产者:测试各种级别
level:0 不延时
1-18分别延时1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
for (int i = 0; i < 10; i++) {
DelayMsg delayMsg = new DelayMsg().setTime(System.currentTimeMillis()).setDelayLevel(i);
Message<DelayMsg> message = MessageBuilder.withPayload(delayMsg).build();
SendResult sendResult = rocketMQTemplate.syncSend("delay_topic", message, 1000, i);
log.info("delayTopicMsg sendResult: {}", sendResult);
}
事务消息
二阶段提交的方式:
首先发送半消息prepared message
指定本地事务,成功则commit message,失败则rollback message,未知状态则回查状态
topic:
首先在dashboard创建事务topic:transaction_topic
消费者
@Component
@Slf4j
@RocketMQMessageListener(topic = "transaction_topic", consumerGroup = "transaction_consumer_group")
public class TransactionConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
log.info("TransactionConsumer onMessage: {}", message);
}
}
本地事务执行及回查
@Slf4j
@RocketMQTransactionListener
public class TransactionListener implements RocketMQLocalTransactionListener {
/**
* 本地事务执行方法
* @param msg 消息
* @param arg 业务参数 监听器中传递额外的上下文信息
*/
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
log.info("executeLocalTransaction: {}, arg: {}", msg.getPayload(), arg);
// 模拟成功、失败、未知逻辑
if ("success".equals(arg)) {
return RocketMQLocalTransactionState.COMMIT;
}
if ("fail".equals(arg)) {
return RocketMQLocalTransactionState.ROLLBACK;
}
return RocketMQLocalTransactionState.UNKNOWN;
}
/**
* executeLocalTransaction 返回 unknown 15s后 回查调用
*/
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
log.info("checkLocalTransaction: {}", msg);
return RocketMQLocalTransactionState.COMMIT;
}
}
生产者:测试成功、失败、未知三种逻辑
List<String> tags = Arrays.asList("success", "fail", "unknown");
for (String tag : tags) {
Message<String> message = MessageBuilder.withPayload("transaction test: " + tag).build();
TransactionSendResult sendResult = rocketMQTemplate.sendMessageInTransaction(
"transaction_topic", message, tag);
log.info("transactionTopicMsg sendResult: {}", sendResult);
}
阿里云ons-client客户端方式
添加依赖
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>ons-client</artifactId>
<!--<version>2.0.8.Final</version>-->
<version>1.8.8.1.Final</version>
</dependency>
2.x版本默认使用gRPC协议,
而本地docker安装的版本不支持gRPC,使用起来有些不方便,所以使用1.8版本
配置消费者接口及实现
public interface RocketMqNormalConsumer extends MessageListener {
String getTopic();
default String getFilterType() {
return "TAG";
}
default String getFilterExpression() {
return "*";
}
}
@Component
@Slf4j
public class NormalConsumerA implements RocketMqNormalConsumer {
@Override
public String getTopic() {
return "normal_topic";
}
@Override
public Action consume(Message message, ConsumeContext consumeContext) {
log.info("NormalConsumerA body {}", new String(message.getBody()));
log.info("NormalConsumerA key {}", message.getKey());
log.info("NormalConsumerA tag {}", message.getTag());
return Action.CommitMessage;
}
}
配置producerBean和consumerBean
private Properties getProperties(String groupId) {
Properties properties = new Properties();
properties.setProperty(PropertyKeyConst.NAMESRV_ADDR, nameServer);
properties.setProperty(PropertyKeyConst.GROUP_ID, groupId);
properties.setProperty(PropertyKeyConst.AccessKey, "not empty");
properties.setProperty(PropertyKeyConst.SecretKey, "not empty");
return properties;
}
@Bean(initMethod = "start", destroyMethod = "shutdown")
public ProducerBean producerBean() {
ProducerBean producerBean = new ProducerBean();
producerBean.setProperties(rocketMqProperties.getProducerProperties());
log.info("producerBean: {}", producerBean);
return producerBean;
}
/**
* 创建一个ConsumerBean,并配置订阅关系
* @param rocketMqNormalConsumers 全部消费者
*/
@Bean(initMethod = "start", destroyMethod = "shutdown")
public ConsumerBean consumerBean(List<RocketMqNormalConsumer> rocketMqNormalConsumers) {
ConsumerBean consumerBean = new ConsumerBean();
consumerBean.setProperties(rocketMqProperties.getConsumerProperties());
Map<Subscription, MessageListener> subscriptionTable = new HashMap<>();
for (RocketMqNormalConsumer rocketMqNormalConsumer : rocketMqNormalConsumers) {
Subscription subscription = new Subscription();
subscription.setTopic(rocketMqNormalConsumer.getTopic());
subscription.setType(rocketMqNormalConsumer.getFilterType());
subscription.setExpression(rocketMqNormalConsumer.getFilterExpression());
subscriptionTable.put(subscription, rocketMqNormalConsumer);
log.info("consumerBean subscription: {} {}", subscription, rocketMqNormalConsumer);
// Subscription [topic=normal_topic, expression=*, type=TAG]
}
consumerBean.setSubscriptionTable(subscriptionTable);
return consumerBean;
}
发送消息
@Resource
private ProducerBean producerBean;
String msg = "normal_topic_msg_" + LocalDateTime.now();
Message message = new Message(topic, "a", "key:" + i, msg.getBytes());
SendResult sendResult = producerBean.send(message);
log.info("sendResult: {}", sendResult);
更多推荐



所有评论(0)