Обучение RocketMQ — публикация сообщений и подписка

задняя часть

В предыдущей статье был проанализирован процесс запуска брокера и рассмотрены основные функции брокера. В следующих нескольких статьях приготовьтесь следоватьНачало работы с RocketMQ за десять минутРяд функций, упомянутых в этой статье, изучаются по очереди. Эта статья готова проанализировать самые основные функции RocketMQ как MQ: публикация сообщений (публикация) и подписка (подписка). Во-первых, я обращаюсь кСерия статей Spring Boot (6): интегрированное использование и мониторинг SpringBoot RocketMQЭтот пост завершает простой пример.

1. Модель сообщений RocketMQ

Скриншот 2018-03-31 14.50.41.png

При развертывании RocketMQ сначала запустите сервер имен, а затем брокер.В это время брокер зарегистрируется на сервере имен. Когда производитель в приложении запускается, он сначала подключается к серверу имен, чтобы получить список адресов брокера, затем он устанавливает соединение с брокером, после чего можно отправлять сообщения. Среди них: один производитель подключается только к одному серверу имен, один производитель устанавливает соединения со всеми брокерами, и каждое соединение будет иметь механизм обнаружения сердцебиения.

Производитель будет опрашивать отправку сообщений в коллекцию mq указанной темы.

У потребителей есть два режима потребления: кластерное потребление и широковещательное потребление. Потребление кластера: несколько потребителей в среднем потребляют все сообщения mq в теме, то есть после того, как сообщение потребляется потребителем в очереди сообщений, другие потребители не будут его потреблять; широковещательное потребление: все потребители могут его потреблять Все сообщения, отправленные на Эта тема.

Потребители могут получать сообщения в двух режимах: push-режим и pull-режим.В RocketMQ, с точки зрения технической реализации, push-режим также является уровнем инкапсуляции в pull-режиме.

2. Отправка сообщения

Демонстрация продюсера

Сначала дайте код,

package com.javadu.chapter8rocketmq.message;

import org.apache.rocketmq.client.exception.MQBrokerException;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import org.apache.rocketmq.remoting.exception.RemotingException;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;

import java.io.UnsupportedEncodingException;

import javax.annotation.PostConstruct;

/**
 * 作用: 同步发送消息
 * User: duqi
 * Date: 2018/3/29
 * Time: 13:52
 */
@Component
public class ProducerDemo {

    @Value("${apache.rocketmq.producer.producerGroup}")
    private String producerGroup;

    @Value("${apache.rocketmq.namesrvAddr}")
    private String namesrvAddr;

    @PostConstruct
    public void defaultMQProducer() {
        DefaultMQProducer defaultMQProducer = new DefaultMQProducer(producerGroup);

        defaultMQProducer.setNamesrvAddr(namesrvAddr);

        try {
            defaultMQProducer.start();

            Message message = new Message("TopicTest", "TagA",
                                          "Hello RocketMQ".getBytes(RemotingHelper.DEFAULT_CHARSET));

            for (int i = 0; i < 100; i++) {
                SendResult sendResult = defaultMQProducer.send(message);
                System.out.println("发送消息结果, msgId:" + sendResult.getMsgId() +
                                   ", 发送状态:" + sendResult.getSendStatus());
            }

        } catch (MQClientException | UnsupportedEncodingException | InterruptedException
            | RemotingException | MQBrokerException e) {
            e.printStackTrace();
        } finally {
            defaultMQProducer.shutdown();
        }
    }

}

У производителя есть два свойства:

  • Адрес сервера имен, который используется для получения информации о брокере
  • В наборе производителей productGroup есть разные экземпляры производителей в одной и той же группе производителей. Если самый ранний производитель выходит из строя, брокер уведомляет другие экземпляры производителей в группе о необходимости фиксации или отката транзакции.

Сообщение в RocketMQ представлено Message, а код определяется следующим образом:

public class Message implements Serializable {
    private static final long serialVersionUID = 8445773977080406428L;

    private String topic;
    private int flag;
    private Map<String, String> properties;
    private byte[] body;

    public Message() {
    }
    //省略了getter和setter方法
}
  • тема: В какую тему будет отправлено сообщение
  • флаг: может использоваться для фильтрации сообщений
  • свойства: расширенное поле, которое может выполнять прозрачную передачу некоторых общих значений системного уровня, таких как segmentId ходьбы по небу
  • тело: содержание сообщения

После отправки каждого сообщения вы получите объект SendResult, посмотрите на структуру объекта:

public class SendResult {
    //发送状态
    private SendStatus sendStatus;
    //消息ID,用于消息去重、消息跟踪
    private String msgId;
    private MessageQueue messageQueue;
    private long queueOffset;
    //事务ID
    private String transactionId;
    private String offsetMsgId;
    private String regionId;
    //是否需要跟踪
    private boolean traceOn = true;

    public SendResult() {
    }
    //省略了构造函数、getter和setter等一系列方法
}

В этой демонстрации мы выводим содержимое сообщения и статус сообщения на консоль вместе.

Анализ исходного кода отправки сообщений

Структура пакета клиентского модуля в RocketMQ выглядит следующим образом.Видно, что автор не поместил определение и реализацию интерфейса в один пакет (это обычная практика в наших бизнес-приложениях, и не обязательно разумная ). Интерфейсы производителя и потребителя определяются в пакетах производителя и потребителя соответственно, а конкретная реализация помещается в пакет impl.

Скриншот 2018-03-31 11.51.36.png

Во-первых, обратите внимание на содержимое пакета производителя.Основные классы следующие: DefaultMQProducer — это реализация производителя по умолчанию, MQAdmin используется для определения некоторых интерфейсов управления, а MQProducer используется для определения некоторых интерфейсов, специфичных для производителя.

MQProducer.png

В ProducerDemo запустите производителя с помощью `defaultMQProducer.start();, а затем посмотрите на процесс метода start():

  • Определите следующее действие на основе статуса службы
  • Для статуса CREATE_JUST
    • установить статус службы
    • Проверить конфигурацию
    • Получите или создайте экземпляр MQClientInstance
    • Зарегистрируйте производителя в указанной группе производителей, то есть в структуре данных таблицы производителей, которая является картой.
    • Заполните структуру данных topicPublishInfoTable.
    • начать продюсер
  • Выдает исключение для RUNNING, START_FAILED и SHUTDOWN_ALREADY
 public void start(final boolean startFactory) throws MQClientException {
        //根据当前的服务状态决定接下来的动作
        switch (this.serviceState) {
            case CREATE_JUST:
                this.serviceState = ServiceState.START_FAILED;

                this.checkConfig();

                if (!this.defaultMQProducer.getProducerGroup().equals(MixAll.CLIENT_INNER_PRODUCER_GROUP)) {
                    this.defaultMQProducer.changeInstanceNameToPID();
                }

                //创建一个客户端工厂
                this.mQClientFactory = MQClientManager.getInstance().getAndCreateMQClientInstance(this.defaultMQProducer, rpcHook);
                //将生产者注册到指定producer group
                boolean registerOK = mQClientFactory.registerProducer(this.defaultMQProducer.getProducerGroup(), this);
                if (!registerOK) {
                    this.serviceState = ServiceState.CREATE_JUST;
                    throw new MQClientException("The producer group[" + this.defaultMQProducer.getProducerGroup()
                        + "] has been created before, specify another name please." + FAQUrl.suggestTodo(FAQUrl.GROUP_NAME_DUPLICATE_URL),
                        null);
                }
                
                //填充topicPublishInfoTable
                this.topicPublishInfoTable.put(this.defaultMQProducer.getCreateTopicKey(), new TopicPublishInfo());

                if (startFactory) {
                    mQClientFactory.start();
                }

                log.info("the producer [{}] start OK. sendMessageWithVIPChannel={}", this.defaultMQProducer.getProducerGroup(),
                    this.defaultMQProducer.isSendMessageWithVIPChannel());
                this.serviceState = ServiceState.RUNNING;
                break;
            case RUNNING:
            case START_FAILED:
            case SHUTDOWN_ALREADY:
                throw new MQClientException("The producer service state not OK, maybe started once, "
                    + this.serviceState
                    + FAQUrl.suggestTodo(FAQUrl.CLIENT_SERVICE_NOT_OK),
                    null);
            default:
                break;
        }

        //给该producer连接的所有broker发送心跳消息
        this.mQClientFactory.sendHeartbeatToAllBrokerWithLock();
    }

следитьmQClientFactory.start()Следуйте ниже, чтобы узнать больше о производителе. Основные шаги:

  • Создайте канал запрос-ответ
  • Запускать различные временные задачи, например: вытягивать адрес кластера брокера на сервер имен каждые 2 минуты, что означает, что, если брокер не работает, доставка сообщения производителя в течение этих двух минут невозможна; периодически вытягивать маршрутную информацию, такую ​​как как темы с сервера имен; регулярно очищайте неисправные брокеры и отправляйте им сообщения пульса.
  • Запустите службу вытягивания, службу балансировки нагрузки, службу выталкивания и другие службы, эти три службы связаны с потребителями. Дизайн здесь не очень четкий, а логика запуска потребителей и производителей сведена воедино. Посмотрите на pullMessageService и rebalanceService и initialization, они инициализируются по MQClientInstance, а MQClientInstance настраивается по ClientConfig.
  public void start() throws MQClientException {

        synchronized (this) {
            switch (this.serviceState) {
                case CREATE_JUST:
                    this.serviceState = ServiceState.START_FAILED;
                    // If not specified,looking address from name server
                    if (null == this.clientConfig.getNamesrvAddr()) {
                        this.mQClientAPIImpl.fetchNameServerAddr();
                    }
                    // Start request-response channel
                    this.mQClientAPIImpl.start();
                    // Start various schedule tasks
                    this.startScheduledTask();
                    // Start pull service
                    this.pullMessageService.start();
                    // Start rebalance service
                    this.rebalanceService.start();
                    // Start push service
                    this.defaultMQProducer.getDefaultMQProducerImpl().start(false);
                    log.info("the client factory [{}] start OK", this.clientId);
                    this.serviceState = ServiceState.RUNNING;
                    break;
                case RUNNING:
                    break;
                case SHUTDOWN_ALREADY:
                    break;
                case START_FAILED:
                    throw new MQClientException("The Factory object[" + this.getClientId() + "] has been created before, and failed.", null);
                default:
                    break;
            }
        }
    }

После того, как производитель запущен, давайте посмотрим на процесс отправки сообщения.Как показано на рисунке ниже, DefaultMQProducer предоставляет множество методов для отправки сообщений, которые могут реализовать синхронную отправку сообщений, асинхронную отправку сообщений, указанную очередь сообщений, одностороннее сообщение, транзакцию. сообщение и т.д.

Скриншот 2018-03-31 12.26.32.png

Здесь мы рассмотрим только самые простыеsend(Message message)метод, который, наконец, реализован в DefaultMQProducerImpl:

    private SendResult sendDefaultImpl(
        Message msg,
        final CommunicationMode communicationMode,
        final SendCallback sendCallback,
        final long timeout
    ) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {
        //确认生产者状态正常
        this.makeSureStateOK();
        //检查消息的合法性
        Validators.checkMessage(msg, this.defaultMQProducer);

        final long invokeID = random.nextLong();
        long beginTimestampFirst = System.currentTimeMillis();
        long beginTimestampPrev = beginTimestampFirst;
        long endTimestamp = beginTimestampFirst;
        //获取消息的目的地:Topic信息
        TopicPublishInfo topicPublishInfo = this.tryToFindTopicPublishInfo(msg.getTopic());
        if (topicPublishInfo != null && topicPublishInfo.ok()) {
            MessageQueue mq = null;
            Exception exception = null;
            SendResult sendResult = null;
            //计算出消息的投递次数,如果是同步投递,则是1+重试次数,如果不是同步投递,则只需要投递一次
            int timesTotal = communicationMode == CommunicationMode.SYNC ? 1 + this.defaultMQProducer.getRetryTimesWhenSendFailed() : 1;
            int times = 0;
            String[] brokersSent = new String[timesTotal];
            //一个broker集群有不同的broker节点,lastBrokerName记录了上次投递的broker节点,每个broker节点
            for (; times < timesTotal; times++) {
                String lastBrokerName = null == mq ? null : mq.getBrokerName();
                //选择一个要发送的消息队列
                MessageQueue mqSelected = this.selectOneMessageQueue(topicPublishInfo, lastBrokerName);
                if (mqSelected != null) {
                    mq = mqSelected;
                    brokersSent[times] = mq.getBrokerName();
                    try {
                        beginTimestampPrev = System.currentTimeMillis();
                        //投递消息
                        sendResult = this.sendKernelImpl(msg, mq, communicationMode, sendCallback, topicPublishInfo, timeout);
                        endTimestamp = System.currentTimeMillis();
                        this.updateFaultItem(mq.getBrokerName(), endTimestamp - beginTimestampPrev, false);
                        //根据消息发送模式,对消息发送结果做不同的处理
                        switch (communicationMode) {
                            case ASYNC:
                                return null;
                            case ONEWAY:
                                return null;
                            case SYNC:
                                if (sendResult.getSendStatus() != SendStatus.SEND_OK) {
                                    if (this.defaultMQProducer.isRetryAnotherBrokerWhenNotStoreOK()) {
                                        continue;
                                    }
                                }

                                return sendResult;
                            default:
                                break;
                        }
                    } catch (RemotingException e) {
                        endTimestamp = System.currentTimeMillis();
                        this.updateFaultItem(mq.getBrokerName(), endTimestamp - beginTimestampPrev, true);
                        log.warn(String.format("sendKernelImpl exception, resend at once, InvokeID: %s, RT: %sms, Broker: %s", invokeID, endTimestamp - beginTimestampPrev, mq), e);
                        log.warn(msg.toString());
                        exception = e;
                        continue;
                    } catch (MQClientException e) {
                        endTimestamp = System.currentTimeMillis();
                        this.updateFaultItem(mq.getBrokerName(), endTimestamp - beginTimestampPrev, true);
                        log.warn(String.format("sendKernelImpl exception, resend at once, InvokeID: %s, RT: %sms, Broker: %s", invokeID, endTimestamp - beginTimestampPrev, mq), e);
                        log.warn(msg.toString());
                        exception = e;
                        continue;
                    } catch (MQBrokerException e) {
                        endTimestamp = System.currentTimeMillis();
                        this.updateFaultItem(mq.getBrokerName(), endTimestamp - beginTimestampPrev, true);
                        log.warn(String.format("sendKernelImpl exception, resend at once, InvokeID: %s, RT: %sms, Broker: %s", invokeID, endTimestamp - beginTimestampPrev, mq), e);
                        log.warn(msg.toString());
                        exception = e;
                        switch (e.getResponseCode()) {
                            case ResponseCode.TOPIC_NOT_EXIST:
                            case ResponseCode.SERVICE_NOT_AVAILABLE:
                            case ResponseCode.SYSTEM_ERROR:
                            case ResponseCode.NO_PERMISSION:
                            case ResponseCode.NO_BUYER_ID:
                            case ResponseCode.NOT_IN_CURRENT_UNIT:
                                continue;
                            default:
                                if (sendResult != null) {
                                    return sendResult;
                                }

                                throw e;
                        }
                    } catch (InterruptedException e) {
                        endTimestamp = System.currentTimeMillis();
                        this.updateFaultItem(mq.getBrokerName(), endTimestamp - beginTimestampPrev, false);
                        log.warn(String.format("sendKernelImpl exception, throw exception, InvokeID: %s, RT: %sms, Broker: %s", invokeID, endTimestamp - beginTimestampPrev, mq), e);
                        log.warn(msg.toString());

                        log.warn("sendKernelImpl exception", e);
                        log.warn(msg.toString());
                        throw e;
                    }
                } else {
                    break;
                }
            }

            if (sendResult != null) {
                return sendResult;
            }

            String info = String.format("Send [%d] times, still failed, cost [%d]ms, Topic: %s, BrokersSent: %s",
                times,
                System.currentTimeMillis() - beginTimestampFirst,
                msg.getTopic(),
                Arrays.toString(brokersSent));

            info += FAQUrl.suggestTodo(FAQUrl.SEND_MSG_FAILED);

            MQClientException mqClientException = new MQClientException(info, exception);
            if (exception instanceof MQBrokerException) {
                mqClientException.setResponseCode(((MQBrokerException) exception).getResponseCode());
            } else if (exception instanceof RemotingConnectException) {
                mqClientException.setResponseCode(ClientErrorCode.CONNECT_BROKER_EXCEPTION);
            } else if (exception instanceof RemotingTimeoutException) {
                mqClientException.setResponseCode(ClientErrorCode.ACCESS_BROKER_TIMEOUT);
            } else if (exception instanceof MQClientException) {
                mqClientException.setResponseCode(ClientErrorCode.BROKER_NOT_EXIST_EXCEPTION);
            }

            throw mqClientException;
        }

        List<String> nsList = this.getmQClientFactory().getMQClientAPIImpl().getNameServerAddressList();
        if (null == nsList || nsList.isEmpty()) {
            throw new MQClientException(
                "No name server address, please set it." + FAQUrl.suggestTodo(FAQUrl.NAME_SERVER_ADDR_NOT_EXIST_URL), null).setResponseCode(ClientErrorCode.NO_NAME_SERVER_EXCEPTION);
        }

        throw new MQClientException("No route info of this topic, " + msg.getTopic() + FAQUrl.suggestTodo(FAQUrl.NO_TOPIC_ROUTE_INFO),
            null).setResponseCode(ClientErrorCode.NOT_FOUND_TOPIC_EXCEPTION);
    }

Основной процесс отправки сообщения выглядит следующим образом:

  • Сначала проверьте легитимность производителя и сообщения
  • Затем получите информацию, отправленную сообщением, которая хранится в объекте TopicPublishInfo:
public class TopicPublishInfo {
    //是否顺序消息
    private boolean orderTopic = false;
    private boolean haveTopicRouterInfo = false;
    //维护该topic下用于的消息队列列表
    private List<MessageQueue> messageQueueList = new ArrayList<MessageQueue>();
    //计算下一次该投递的队列,这里应用ThreadLocal,即使是同一台机器中,每个producer实例都有自己的队列
    private volatile ThreadLocalIndex sendWhichQueue = new ThreadLocalIndex();
    private TopicRouteData topicRouteData;

    //省略了getter和setter方法
    
    //选择指定lastBrokerName上的下一个mq
    public MessageQueue selectOneMessageQueue(final String lastBrokerName) {
        if (lastBrokerName == null) {
            return selectOneMessageQueue();
        } else {
            int index = this.sendWhichQueue.getAndIncrement();
            for (int i = 0; i < this.messageQueueList.size(); i++) {
                int pos = Math.abs(index++) % this.messageQueueList.size();
                if (pos < 0)
                    pos = 0;
                MessageQueue mq = this.messageQueueList.get(pos);
                if (!mq.getBrokerName().equals(lastBrokerName)) {
                    return mq;
                }
            }
            return selectOneMessageQueue();
        }
    }

    //选择当前broker节点的下一个mq
    public MessageQueue selectOneMessageQueue() {
        int index = this.sendWhichQueue.getAndIncrement();
        int pos = Math.abs(index) % this.messageQueueList.size();
        if (pos < 0)
            pos = 0;
        return this.messageQueueList.get(pos);
    }
}
  • Выберите MessageQueue для отправки в топик.Логика выбора делится на два случая: (1) По умолчанию на узле брокера, который был доставлен в последний раз, опрашивает следующую очередь сообщений для отправки, (2) Значение of sendLatencyFaultEnable При значении true этот кусок не очень понятен.
  • опубликовать сообщение
  • В зависимости от режима работы очереди сообщений выполняется различная обработка результата доставки.

2. Потребление сообщений

Потребительская демонстрация

В потребителе есть атрибут, на который нужно обратить внимание:

  • ConsumerGroup: Экземпляры-потребители в той же ConsumerGroup имеют те же роли, что и экземпляры-производители в productGroup; экземпляры в ConsumerGroup также могут реализовывать балансировку нагрузки и аварийное восстановление. PS: Экземпляры потребителей в одной и той же потребительской группе должны быть подписаны на одну и ту же тему.
  • Address of nameServer: адрес сервера имен, используемый для получения информации о брокере и топике.

В потребительской демонстрации выполняются следующие действия:

  • Установить свойства конфигурации
  • Задайте тему подписки, можно указать тег
  • При настройке первого запуска откуда начинать потребление из очереди сообщений
  • установить обработчик сообщений
  • начать потребитель
package com.javadu.chapter8rocketmq.message;

import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;

import javax.annotation.PostConstruct;

/**
 * 作用:
 * User: duqi
 * Date: 2018/3/29
 * Time: 14:00
 */
@Component
public class ConsumerDemo {

    /**
     * 消费者的组名
     */
    @Value("${apache.rocketmq.consumer.consumerGroup}")
    private String consumerGroup;

    /**
     * NameServer 地址
     */
    @Value("${apache.rocketmq.namesrvAddr}")
    private String namesrvAddr;

    @PostConstruct
    public void defaultMQPushConsumer() {
        //消费者的组名
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(consumerGroup);

        //指定NameServer地址,多个地址以 ; 隔开
        consumer.setNamesrvAddr(namesrvAddr);
        try {
            //订阅PushTopic下Tag为push的消息
            consumer.subscribe("TopicTest", "TagA");

            //设置Consumer第一次启动是从队列头部开始消费还是队列尾部开始消费
            //如果非第一次启动,那么按照上次消费的位置继续消费
            consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
            consumer.registerMessageListener((MessageListenerConcurrently) (list, context) -> {
                try {
                    for (MessageExt messageExt : list) {

                        //输出消息内容
                        System.out.println("messageExt: " + messageExt);

                        String messageBody = new String(messageExt.getBody(), RemotingHelper.DEFAULT_CHARSET);

                        //输出消息内容
                        System.out.println("消费响应:msgId : " + messageExt.getMsgId() + ",  msgBody : " + messageBody);
                    }
                } catch (Exception e) {
                    e.printStackTrace();
                    //稍后再试
                    return ConsumeConcurrentlyStatus.RECONSUME_LATER;
                }
                //消费成功
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            });
            consumer.start();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

Анализ исходного кода потребителя

Как было проанализировано ранее, клиентский модуль в RocketMQ унифицированным образом предоставляет как производителей, так и потребителей.Давайте рассмотрим несколько основных классов потребителей. Как упоминалось ранее, RocketMQ на самом деле является режимом извлечения.. DefaultMQPushConsumer здесь реализует режим push и только инкапсулирует службу сообщений извлечения, то есть, когда сообщение извлекается, бизнес-потребитель инициируется для регистрации обратного вызова здесь, и конкретный Служба pull-сообщений реализована в PullMessageService, и эта деталь будет рассмотрена позже.

MQConsumer.png

В ConsumerDemo после установки конфигурационной информации будет выполнена подписка на тему и будет вызван метод подписки DefaultMQPushConsumer.Исходный код выглядит следующим образом:

    /**
     * Subscribe a topic to consuming subscription.
     *
     * @param topic topic to subscribe.
     * @param subExpression subscription expression.it only support or operation such as "tag1 || tag2 || tag3" <br>
     * if null or * expression,meaning subscribe all
     * @throws MQClientException if there is any client error.
     */
    @Override
    public void subscribe(String topic, String subExpression) throws MQClientException {
        this.defaultMQPushConsumerImpl.subscribe(topic, subExpression);
    }

Первый параметр — это информация о теме, а второй параметр используется для поля тега фильтрации сообщений. Настоящая подписка происходит в DefaultMQPushConsumerImpl, код такой:

    public void subscribe(String topic, String subExpression) throws MQClientException {
        try {
            //构建包含订阅信息的对象,并放入负载平衡组件维护的map中,以topic为key
            SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData(this.defaultMQPushConsumer.getConsumerGroup(),
                topic, subExpression);
            this.rebalanceImpl.getSubscriptionInner().put(topic, subscriptionData);
            //如果已经跟broker集群建立连接,则给所有的broker节点发送心跳消息
            if (this.mQClientFactory != null) {
                this.mQClientFactory.sendHeartbeatToAllBrokerWithLock();
            }
        } catch (Exception e) {
            throw new MQClientException("subscription exception", e);
        }
    }

Следующим шагом в ConsumerDemo является установка начальной позиции сообщения потребителя при первом запуске потребителя.Для этого используется атрибут в DefaultMQPushConsumer — ConsumerFromWhere, который имеет три возможных значения.

  • CONSUME_FROM_LAST_OFFSET, значение по умолчанию, означающее начать потребление с того места, где оно было остановлено в последний раз
  • CONSUME_FROM_FIRST_OFFSET, начать потребление с начала очереди
  • CONSUME_FROM_TIMESTAMP, начать потребление с указанного момента времени

Затем ConsumerDemo зарегистрирует обратный вызов и обработает сообщение, когда оно прибудет (последний прослушиватель сообщений поддерживает параллельное потребление):

    /**
     * Register a callback to execute on message arrival for concurrent consuming.
     *
     * @param messageListener message handling callback.
     */
    @Override
    public void registerMessageListener(MessageListenerConcurrently messageListener) {
        this.messageListener = messageListener;
        this.defaultMQPushConsumerImpl.registerMessageListener(messageListener);
    }

Наконец, давайте взглянем на процесс запуска ConsumerDemo, то есть метод запуска DefaultMQPushConsumerImpl, который в основном выполняет следующие действия:

  • Проверить конфигурацию
  • Скопируйте информацию о подписке в компонент балансировки нагрузки (rebalanceImpl);
  • Настройка нескольких свойств компонента балансировки нагрузки
  • Обработка конфигурации различных режимов сообщений (кластерный режим или широковещательный режим)
  • Различные конфигурации для обработки последовательного и одновременного потребления
  • Зарегистрируйте информацию о потребителе и группу потребителей в ConsumerTable экземпляра клиента MQ.
  • Запустите потребительский клиент

использованная литература

  1. Принцип и практика распределенной открытой системы обмена сообщениями (RocketMQ)
  2. Rocketmq-spring-boot-starter при условии покупки хорошей машины
  3. Серия статей Spring Boot (6): интегрированное использование и мониторинг SpringBoot RocketMQ