暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

RocketMQ源码分析之发消息

徘徊笔记 2019-04-13
327

来源:https://github.com/apache/rocketmq


生产者发送消息,消息发送默认时间为3000ms,默认是同步发送

    producer.send(msg);
    publicSendResult send(
    Message msg) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {
    return this.defaultMQProducerImpl.send(msg);
    }
    publicSendResult send(
    Message msg) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {
    return send(msg, this.defaultMQProducer.getSendMsgTimeout());
    }
    public SendResult send(Message msg,
    long timeout) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {
    return this.sendDefaultImpl(msg, CommunicationMode.SYNC, null, timeout);
    }


    校验topic长度是否合理不为空不能为"TBW102"并且长度不能超过255消息

      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;
      TopicPublishInfo topicPublishInfo = this.tryToFindTopicPublishInfo(msg.getTopic());
      if (topicPublishInfo != null && topicPublishInfo.ok()) {
      boolean callTimeout = false;
      MessageQueue mq = null;
      Exception exception = null;
      SendResult sendResult = null;
      int timesTotal = communicationMode == CommunicationMode.SYNC ? 1 + this.defaultMQProducer.getRetryTimesWhenSendFailed() : 1;
      int times = 0;
      String[] brokersSent = new String[timesTotal];
      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();
      long costTime = beginTimestampPrev - beginTimestampFirst;
      if (timeout < costTime) {
      callTimeout = true;
      break;
      }


      sendResult = this.sendKernelImpl(msg, mq, communicationMode, sendCallback, topicPublishInfo, timeout - costTime);
      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 (callTimeout) {
      throw new RemotingTooMuchRequestException("sendDefaultImpl call timeout");
      }


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


      查找消息的topic所对应的broker发布信息,先查找有没有,没有的话就需要从配置中心进行拉取,继续判断有没有topic的路由信息以及对应的消息队列是否存在,当然如果是自动创建topic的话,也就是autoCreateTopicEnable = true时,第一次从nameserver是查不到该topic信息的,因为所有的broker是没有该topic信息的,所以第二次会进行默认topic查询,因为当该标识为true时,broker会创建默认topic,所以也就会查到对应的broker消息队列信息返回。当然如果设置autoCreateTopicEnable = false,并且没有在broker配置该topic,就会报错"No route info of this topic, "后面会随机挑选一个主broker进行发送消息,所以以后这个topic只会存在这一个主broker里,当然该broker的从节点上还是可以读取消息的,其他的主broker是没有该topic,所以如果你想让每个topic都存在每一个主broker,建议设置该值为false,一个是有利于统一管理,另一个你可以发消息给所有配置有该topic的主broker。

        private TopicPublishInfo tryToFindTopicPublishInfo(final String topic) {
        TopicPublishInfo topicPublishInfo = this.topicPublishInfoTable.get(topic);
        if (null == topicPublishInfo || !topicPublishInfo.ok()) {
        this.topicPublishInfoTable.putIfAbsent(topic, new TopicPublishInfo());
        this.mQClientFactory.updateTopicRouteInfoFromNameServer(topic);
        topicPublishInfo = this.topicPublishInfoTable.get(topic);
        }


        if (topicPublishInfo.isHaveTopicRouterInfo() || topicPublishInfo.ok()) {
        return topicPublishInfo;
        } else {
        this.mQClientFactory.updateTopicRouteInfoFromNameServer(topic, true, this.defaultMQProducer);
        topicPublishInfo = this.topicPublishInfoTable.get(topic);
        return topicPublishInfo;
        }
        }


        同步发送默认重试2次,所以如果有异常的话就会总共发送三次,从broker中根绝选择算法来选择一个消息队列,当然重试的话就尽量不会选择之前选择过的broker,当然如果只有这一个主broker的话,三次重试都是发到该地址。

          publicMessageQueue selectOneMessageQueue(final TopicPublishInfo tpInfo, final String lastBrokerName) {
          return this.mqFaultStrategy.selectOneMessageQueue(tpInfo, lastBrokerName);
          }
          public MessageQueue selectOneMessageQueue(final TopicPublishInfo tpInfo, final String lastBrokerName) {
          if (this.sendLatencyFaultEnable) {
          try {
          int index = tpInfo.getSendWhichQueue().getAndIncrement();
          for (int i = 0; i < tpInfo.getMessageQueueList().size(); i++) {
          int pos = Math.abs(index++) % tpInfo.getMessageQueueList().size();
          if (pos < 0)
          pos = 0;
          MessageQueue mq = tpInfo.getMessageQueueList().get(pos);
          if (latencyFaultTolerance.isAvailable(mq.getBrokerName())) {
          if (null == lastBrokerName || mq.getBrokerName().equals(lastBrokerName))
          return mq;
          }
          }


          final String notBestBroker = latencyFaultTolerance.pickOneAtLeast();
          int writeQueueNums = tpInfo.getQueueIdByBroker(notBestBroker);
          if (writeQueueNums > 0) {
          final MessageQueue mq = tpInfo.selectOneMessageQueue();
          if (notBestBroker != null) {
          mq.setBrokerName(notBestBroker);
          mq.setQueueId(tpInfo.getSendWhichQueue().getAndIncrement() % writeQueueNums);
          }
          return mq;
          } else {
          latencyFaultTolerance.remove(notBestBroker);
          }
          } catch (Exception e) {
          log.error("Error occurred when selecting message queue", e);
          }


          return tpInfo.selectOneMessageQueue();
          }


          return tpInfo.selectOneMessageQueue(lastBrokerName);
          }
          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();
          }
          }


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


          获取主broker的地址信息,如果不存在的话就重新从配置中心更新一下topic配置等信息,根据vipChannelEnabled配置来选择与broker的普通服务端通信还是快速响应服务端通讯,因为两者的端口号相差2,给单个消息生成唯一key,批量消息在消息组合成MessageBatch时已经设置过了

            public String findBrokerAddressInPublish(final String brokerName) {
            HashMap<Long/* brokerId */, String/* address */> map = this.brokerAddrTable.get(brokerName);
            if (map != null && !map.isEmpty()) {
            return map.get(MixAll.MASTER_ID);
                }
            return null;
            }
            private SendResult sendKernelImpl(final Message msg,
            final MessageQueue mq,
            final CommunicationMode communicationMode,
            final SendCallback sendCallback,
            final TopicPublishInfo topicPublishInfo,
            final long timeout) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {
            long beginStartTime = System.currentTimeMillis();
            String brokerAddr = this.mQClientFactory.findBrokerAddressInPublish(mq.getBrokerName());
            if (null == brokerAddr) {
            tryToFindTopicPublishInfo(mq.getTopic());
            brokerAddr = this.mQClientFactory.findBrokerAddressInPublish(mq.getBrokerName());
            }


            SendMessageContext context = null;
            if (brokerAddr != null) {
            brokerAddr = MixAll.brokerVIPChannel(this.defaultMQProducer.isSendMessageWithVIPChannel(), brokerAddr);


            byte[] prevBody = msg.getBody();
            try {
            //for MessageBatch,ID has been set in the generating process
            if (!(msg instanceof MessageBatch)) {
            MessageClientIDSetter.setUniqID(msg);
            }


            int sysFlag = 0;
            boolean msgBodyCompressed = false;
            if (this.tryToCompressMessage(msg)) {
            sysFlag |= MessageSysFlag.COMPRESSED_FLAG;
            msgBodyCompressed = true;
            }


            final String tranMsg = msg.getProperty(MessageConst.PROPERTY_TRANSACTION_PREPARED);
            if (tranMsg != null && Boolean.parseBoolean(tranMsg)) {
            sysFlag |= MessageSysFlag.TRANSACTION_PREPARED_TYPE;
            }


            if (hasCheckForbiddenHook()) {
            CheckForbiddenContext checkForbiddenContext = new CheckForbiddenContext();
            checkForbiddenContext.setNameSrvAddr(this.defaultMQProducer.getNamesrvAddr());
            checkForbiddenContext.setGroup(this.defaultMQProducer.getProducerGroup());
            checkForbiddenContext.setCommunicationMode(communicationMode);
            checkForbiddenContext.setBrokerAddr(brokerAddr);
            checkForbiddenContext.setMessage(msg);
            checkForbiddenContext.setMq(mq);
            checkForbiddenContext.setUnitMode(this.isUnitMode());
            this.executeCheckForbiddenHook(checkForbiddenContext);
            }


            if (this.hasSendMessageHook()) {
            context = new SendMessageContext();
            context.setProducer(this);
            context.setProducerGroup(this.defaultMQProducer.getProducerGroup());
            context.setCommunicationMode(communicationMode);
            context.setBornHost(this.defaultMQProducer.getClientIP());
            context.setBrokerAddr(brokerAddr);
            context.setMessage(msg);
            context.setMq(mq);
            String isTrans = msg.getProperty(MessageConst.PROPERTY_TRANSACTION_PREPARED);
            if (isTrans != null && isTrans.equals("true")) {
            context.setMsgType(MessageType.Trans_Msg_Half);
            }


            if (msg.getProperty("__STARTDELIVERTIME") != null || msg.getProperty(MessageConst.PROPERTY_DELAY_TIME_LEVEL) != null) {
            context.setMsgType(MessageType.Delay_Msg);
            }
            this.executeSendMessageHookBefore(context);
            }


            SendMessageRequestHeader requestHeader = new SendMessageRequestHeader();
            requestHeader.setProducerGroup(this.defaultMQProducer.getProducerGroup());
            requestHeader.setTopic(msg.getTopic());
            requestHeader.setDefaultTopic(this.defaultMQProducer.getCreateTopicKey());
            requestHeader.setDefaultTopicQueueNums(this.defaultMQProducer.getDefaultTopicQueueNums());
            requestHeader.setQueueId(mq.getQueueId());
            requestHeader.setSysFlag(sysFlag);
            requestHeader.setBornTimestamp(System.currentTimeMillis());
            requestHeader.setFlag(msg.getFlag());
            requestHeader.setProperties(MessageDecoder.messageProperties2String(msg.getProperties()));
            requestHeader.setReconsumeTimes(0);
            requestHeader.setUnitMode(this.isUnitMode());
            requestHeader.setBatch(msg instanceof MessageBatch);
            if (requestHeader.getTopic().startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
            String reconsumeTimes = MessageAccessor.getReconsumeTime(msg);
            if (reconsumeTimes != null) {
            requestHeader.setReconsumeTimes(Integer.valueOf(reconsumeTimes));
            MessageAccessor.clearProperty(msg, MessageConst.PROPERTY_RECONSUME_TIME);
            }


            String maxReconsumeTimes = MessageAccessor.getMaxReconsumeTimes(msg);
            if (maxReconsumeTimes != null) {
            requestHeader.setMaxReconsumeTimes(Integer.valueOf(maxReconsumeTimes));
            MessageAccessor.clearProperty(msg, MessageConst.PROPERTY_MAX_RECONSUME_TIMES);
            }
            }


            SendResult sendResult = null;
            switch (communicationMode) {
            case ASYNC:
            Message tmpMessage = msg;
            if (msgBodyCompressed) {
            //If msg body was compressed, msgbody should be reset using prevBody.
            //Clone new message using commpressed message body and recover origin massage.
            //Fix bug:https://github.com/apache/rocketmq-externals/issues/66
            tmpMessage = MessageAccessor.cloneMessage(msg);
            msg.setBody(prevBody);
            }
            long costTimeAsync = System.currentTimeMillis() - beginStartTime;
            if (timeout < costTimeAsync) {
            throw new RemotingTooMuchRequestException("sendKernelImpl call timeout");
            }
            sendResult = this.mQClientFactory.getMQClientAPIImpl().sendMessage(
            brokerAddr,
            mq.getBrokerName(),
            tmpMessage,
            requestHeader,
            timeout - costTimeAsync,
            communicationMode,
            sendCallback,
            topicPublishInfo,
            this.mQClientFactory,
            this.defaultMQProducer.getRetryTimesWhenSendAsyncFailed(),
            context,
            this);
            break;
            case ONEWAY:
            case SYNC:
            long costTimeSync = System.currentTimeMillis() - beginStartTime;
            if (timeout < costTimeSync) {
            throw new RemotingTooMuchRequestException("sendKernelImpl call timeout");
            }
            sendResult = this.mQClientFactory.getMQClientAPIImpl().sendMessage(
            brokerAddr,
            mq.getBrokerName(),
            msg,
            requestHeader,
            timeout - costTimeSync,
            communicationMode,
            context,
            this);
            break;
            default:
            assert false;
            break;
            }


            if (this.hasSendMessageHook()) {
            context.setSendResult(sendResult);
            this.executeSendMessageHookAfter(context);
            }


            return sendResult;
            } catch (RemotingException e) {
            if (this.hasSendMessageHook()) {
            context.setException(e);
            this.executeSendMessageHookAfter(context);
            }
            throw e;
            } catch (MQBrokerException e) {
            if (this.hasSendMessageHook()) {
            context.setException(e);
            this.executeSendMessageHookAfter(context);
            }
            throw e;
            } catch (InterruptedException e) {
            if (this.hasSendMessageHook()) {
            context.setException(e);
            this.executeSendMessageHookAfter(context);
            }
            throw e;
            } finally {
            msg.setBody(prevBody);
            }
            }


            throw new MQClientException("The broker[" + mq.getBrokerName() + "] not exist", null);
            }


            批量消息不支持压缩,单个消息的话默认超过4K就要进行压缩,判断是否是事务消息,设置消息flag标志,标志位如下,如果是重试消息的话就需要重新设置重试次数,如果是异步发送消息的话就需要防止压缩后改变原消息的body,因为异步操作可能会用到原消息

              public final static int COMPRESSED_FLAG = 0x1;
              public final static int MULTI_TAGS_FLAG = 0x1 << 1;
              public final static int TRANSACTION_NOT_TYPE = 0;
              public final static int TRANSACTION_PREPARED_TYPE = 0x1 << 2;
              public final static int TRANSACTION_COMMIT_TYPE = 0x2 << 2;
              public final static int TRANSACTION_ROLLBACK_TYPE = 0x3 << 2;


              发送消息,判断消息类型组装不同的消息头,使用短变量名加速fastjson反序列化过程。

                public SendResult sendMessage(
                final String addr,
                final String brokerName,
                final Message msg,
                final SendMessageRequestHeader requestHeader,
                final long timeoutMillis,
                final CommunicationMode communicationMode,
                final SendCallback sendCallback,
                final TopicPublishInfo topicPublishInfo,
                final MQClientInstance instance,
                final int retryTimesWhenSendFailed,
                final SendMessageContext context,
                final DefaultMQProducerImpl producer
                )throws RemotingException, MQBrokerException, InterruptedException {
                long beginStartTime = System.currentTimeMillis();
                RemotingCommand request = null;
                if (sendSmartMsg || msg instanceof MessageBatch) {
                SendMessageRequestHeaderV2 requestHeaderV2 = SendMessageRequestHeaderV2.createSendMessageRequestHeaderV2(requestHeader);
                request = RemotingCommand.createRequestCommand(msg instanceof MessageBatch ? RequestCode.SEND_BATCH_MESSAGE : RequestCode.SEND_MESSAGE_V2, requestHeaderV2);
                } else {
                request = RemotingCommand.createRequestCommand(RequestCode.SEND_MESSAGE, requestHeader);
                }


                request.setBody(msg.getBody());


                switch (communicationMode) {
                case ONEWAY:
                this.remotingClient.invokeOneway(addr, request, timeoutMillis);
                return null;
                case ASYNC:
                final AtomicInteger times = new AtomicInteger();
                long costTimeAsync = System.currentTimeMillis() - beginStartTime;
                if (timeoutMillis < costTimeAsync) {
                throw new RemotingTooMuchRequestException("sendMessage call timeout");
                }
                this.sendMessageAsync(addr, brokerName, msg, timeoutMillis - costTimeAsync, request, sendCallback, topicPublishInfo, instance,
                retryTimesWhenSendFailed, times, context, producer);
                return null;
                case SYNC:
                long costTimeSync = System.currentTimeMillis() - beginStartTime;
                if (timeoutMillis < costTimeSync) {
                throw new RemotingTooMuchRequestException("sendMessage call timeout");
                }
                return this.sendMessageSync(addr, brokerName, msg, timeoutMillis - costTimeSync, request);
                default:
                assert false;
                break;
                }


                return null;
                }


                同步发送消息

                  RemotingCommand response = this.remotingClient.invokeSync(addr, request, timeoutMillis);


                  文章转载自徘徊笔记,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

                  评论