回顾一下springboot集成rocketmq的一些用法,实例进行测试。

目录

rocketmq简介

基础概念

消息的发送与消费

消息的可靠性和高可用性

docker安装rocketmq

启动nameserver

启动broker

安装dashboard

springboot集成rocketmq

原生rocketmq客户端方式

简单Demo

普通消息

顺序消息

延时/定时消息

事务消息

阿里云ons-client客户端方式


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);

更多推荐