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

RocketMQ源码分析之消费消息(上)

徘徊笔记 2019-05-04
463

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


当从broker拉取到消息后,开始执行回调方法

    private void executeInvokeCallback(final ResponseFuture responseFuture) {
    boolean runInThisThread = false;
    ExecutorService executor = this.getCallbackExecutor();
    if (executor != null) {
    try {
    executor.submit(new Runnable() {
    @Override
    public void run() {
    try {
    responseFuture.executeInvokeCallback();
    } catch (Throwable e) {
    log.warn("execute callback in executor exception, and callback throw", e);
    } finally {
    responseFuture.release();
    }
    }
    });
    } catch (Exception e) {
    runInThisThread = true;
    log.warn("execute callback in executor exception, maybe executor busy", e);
    }
    } else {
    runInThisThread = true;
    }


    if (runInThisThread) {
    try {
    responseFuture.executeInvokeCallback();
    } catch (Throwable e) {
    log.warn("executeInvokeCallback Exception", e);
    } finally {
    responseFuture.release();
    }
    }
    }


    并且回调方法只能执行一次,解析具体回应的数据

      public void executeInvokeCallback() {
      if (invokeCallback != null) {
      if (this.executeCallbackOnlyOnce.compareAndSet(false, true)) {
      invokeCallback.operationComplete(this);
      }
      }
      }


      public void operationComplete(ResponseFuture responseFuture) {
      RemotingCommand response = responseFuture.getResponseCommand();
      if (response != null) {
      try {
      PullResult pullResult = MQClientAPIImpl.this.processPullResponse(response);
      assert pullResult != null;
      pullCallback.onSuccess(pullResult);
      } catch (Exception e) {
      pullCallback.onException(e);
      }
      } else {
      if (!responseFuture.isSendRequestOK()) {
      pullCallback.onException(new MQClientException("send request failed to " + addr + ". Request: " + request, responseFuture.getCause()));
      } else if (responseFuture.isTimeout()) {
      pullCallback.onException(new MQClientException("wait response from " + addr + " timeout :" + responseFuture.getTimeoutMillis() + "ms" + ". Request: " + request,
      responseFuture.getCause()));
      } else {
      pullCallback.onException(new MQClientException("unknown reason. addr: " + addr + ", timeoutMillis: " + timeoutMillis + ". Request: " + request, responseFuture.getCause()));
      }
      }
      }


      private PullResult processPullResponse(final RemotingCommand response
      throws MQBrokerException, RemotingCommandException {
      PullStatus pullStatus = PullStatus.NO_NEW_MSG;
      switch (response.getCode()) {
      case ResponseCode.SUCCESS:
      pullStatus = PullStatus.FOUND;
      break;
      case ResponseCode.PULL_NOT_FOUND:
      pullStatus = PullStatus.NO_NEW_MSG;
      break;
      case ResponseCode.PULL_RETRY_IMMEDIATELY:
      pullStatus = PullStatus.NO_MATCHED_MSG;
      break;
      case ResponseCode.PULL_OFFSET_MOVED:
      pullStatus = PullStatus.OFFSET_ILLEGAL;
      break;


      default:
      throw new MQBrokerException(response.getCode(), response.getRemark());
      }


      PullMessageResponseHeader responseHeader =
      (PullMessageResponseHeader) response.decodeCommandCustomHeader(PullMessageResponseHeader.class);


      return new PullResultExt(pullStatus, responseHeader.getNextBeginOffset(), responseHeader.getMinOffset(),
      responseHeader.getMaxOffset(), null, responseHeader.getSuggestWhichBrokerId(), response.getBody());
      }


      当解析出现异常时设置待会再进行拉取。

        PullCallback pullCallback = new PullCallback() {
        @Override
        public void onSuccess(PullResult pullResult) {
        if (pullResult != null) {
        pullResult = DefaultMQPushConsumerImpl.this.pullAPIWrapper.processPullResult(pullRequest.getMessageQueue(), pullResult,
        subscriptionData);


        switch (pullResult.getPullStatus()) {
        case FOUND:
        long prevRequestOffset = pullRequest.getNextOffset();
        pullRequest.setNextOffset(pullResult.getNextBeginOffset());
        long pullRT = System.currentTimeMillis() - beginTimestamp;
        DefaultMQPushConsumerImpl.this.getConsumerStatsManager().incPullRT(pullRequest.getConsumerGroup(),
        pullRequest.getMessageQueue().getTopic(), pullRT);


        long firstMsgOffset = Long.MAX_VALUE;
        if (pullResult.getMsgFoundList() == null || pullResult.getMsgFoundList().isEmpty()) {
        DefaultMQPushConsumerImpl.this.executePullRequestImmediately(pullRequest);
        } else {
        firstMsgOffset = pullResult.getMsgFoundList().get(0).getQueueOffset();


        DefaultMQPushConsumerImpl.this.getConsumerStatsManager().incPullTPS(pullRequest.getConsumerGroup(),
        pullRequest.getMessageQueue().getTopic(), pullResult.getMsgFoundList().size());


        boolean dispatchToConsume = processQueue.putMessage(pullResult.getMsgFoundList());
        DefaultMQPushConsumerImpl.this.consumeMessageService.submitConsumeRequest(
        pullResult.getMsgFoundList(),
        processQueue,
        pullRequest.getMessageQueue(),
        dispatchToConsume);


        if (DefaultMQPushConsumerImpl.this.defaultMQPushConsumer.getPullInterval() > 0) {
        DefaultMQPushConsumerImpl.this.executePullRequestLater(pullRequest,
        DefaultMQPushConsumerImpl.this.defaultMQPushConsumer.getPullInterval());
        } else {
        DefaultMQPushConsumerImpl.this.executePullRequestImmediately(pullRequest);
        }
        }


        if (pullResult.getNextBeginOffset() < prevRequestOffset
        || firstMsgOffset < prevRequestOffset) {
        log.warn(
        "[BUG] pull message result maybe data wrong, nextBeginOffset: {} firstMsgOffset: {} prevRequestOffset: {}",
        pullResult.getNextBeginOffset(),
        firstMsgOffset,
        prevRequestOffset);
        }


        break;
        case NO_NEW_MSG:
        pullRequest.setNextOffset(pullResult.getNextBeginOffset());


        DefaultMQPushConsumerImpl.this.correctTagsOffset(pullRequest);


        DefaultMQPushConsumerImpl.this.executePullRequestImmediately(pullRequest);
        break;
        case NO_MATCHED_MSG:
        pullRequest.setNextOffset(pullResult.getNextBeginOffset());


        DefaultMQPushConsumerImpl.this.correctTagsOffset(pullRequest);


        DefaultMQPushConsumerImpl.this.executePullRequestImmediately(pullRequest);
        break;
        case OFFSET_ILLEGAL:
        log.warn("the pull request offset illegal, {} {}",
        pullRequest.toString(), pullResult.toString());
        pullRequest.setNextOffset(pullResult.getNextBeginOffset());


        pullRequest.getProcessQueue().setDropped(true);
        DefaultMQPushConsumerImpl.this.executeTaskLater(new Runnable() {


        @Override
        public void run() {
        try {
        DefaultMQPushConsumerImpl.this.offsetStore.updateOffset(pullRequest.getMessageQueue(),
        pullRequest.getNextOffset(), false);


        DefaultMQPushConsumerImpl.this.offsetStore.persist(pullRequest.getMessageQueue());


        DefaultMQPushConsumerImpl.this.rebalanceImpl.removeProcessQueue(pullRequest.getMessageQueue());


        log.warn("fix the pull request offset, {}", pullRequest);
        } catch (Throwable e) {
        log.error("executeTaskLater Exception", e);
        }
        }
        }, 10000);
        break;
        default:
        break;
        }
        }
        }


        @Override
        public void onException(Throwable e) {
        if (!pullRequest.getMessageQueue().getTopic().startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
        log.warn("execute the pull request exception", e);
        }


        DefaultMQPushConsumerImpl.this.executePullRequestLater(pullRequest, PULL_TIME_DELAY_MILLS_WHEN_EXCEPTION);
        }
        };


        对拉取结果进行解析,更新下一次拉取的broker节点,当有新消息时,对消息进行解码,然后与订阅的tag进行对比过滤,这次不是对比哈希值,而是直接对比tag值,防止哈希冲突造成的错误,给消息设置最大最小偏移量事务id属性

          public void updatePullFromWhichNode(final MessageQueue mq, final long brokerId) {
          AtomicLong suggest = this.pullFromWhichNodeTable.get(mq);
          if (null == suggest) {
          this.pullFromWhichNodeTable.put(mq, new AtomicLong(brokerId));
          } else {
          suggest.set(brokerId);
          }
          }


          public PullResult processPullResult(final MessageQueue mq, final PullResult pullResult,
          final SubscriptionData subscriptionData) {
          PullResultExt pullResultExt = (PullResultExt) pullResult;


          this.updatePullFromWhichNode(mq, pullResultExt.getSuggestWhichBrokerId());
          if (PullStatus.FOUND == pullResult.getPullStatus()) {
          ByteBuffer byteBuffer = ByteBuffer.wrap(pullResultExt.getMessageBinary());
          List<MessageExt> msgList = MessageDecoder.decodes(byteBuffer);


          List<MessageExt> msgListFilterAgain = msgList;
          if (!subscriptionData.getTagsSet().isEmpty() && !subscriptionData.isClassFilterMode()) {
          msgListFilterAgain = new ArrayList<MessageExt>(msgList.size());
          for (MessageExt msg : msgList) {
          if (msg.getTags() != null) {
          if (subscriptionData.getTagsSet().contains(msg.getTags())) {
          msgListFilterAgain.add(msg);
          }
          }
          }
          }


          if (this.hasHook()) {
          FilterMessageContext filterMessageContext = new FilterMessageContext();
          filterMessageContext.setUnitMode(unitMode);
          filterMessageContext.setMsgList(msgListFilterAgain);
          this.executeHook(filterMessageContext);
          }


          for (MessageExt msg : msgListFilterAgain) {
          String traFlag = msg.getProperty(MessageConst.PROPERTY_TRANSACTION_PREPARED);
          if (traFlag != null && Boolean.parseBoolean(traFlag)) {
          msg.setTransactionId(msg.getProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX));
          }
          MessageAccessor.putProperty(msg, MessageConst.PROPERTY_MIN_OFFSET,
          Long.toString(pullResult.getMinOffset()));
          MessageAccessor.putProperty(msg, MessageConst.PROPERTY_MAX_OFFSET,
          Long.toString(pullResult.getMaxOffset()));
          }


          pullResultExt.setMsgFoundList(msgListFilterAgain);
          }


          pullResultExt.setMessageBinary(null);


          return pullResult;
          }


          当没有新消息时,更新本地保存的队列偏移量,然后立刻开始下一次拉取。当发现新消息时,把消息放进处理队列中,记录有效的消息数量,记录消息大小,判断消息是否在被消费,记录最后一个消息与broker的topic对应的队列的下标的跨度,提交消费请求,根据间隔时间判断是否快速启动下一次请求。

            public boolean putMessage(final List<MessageExt> msgs) {
            boolean dispatchToConsume = false;
            try {
            this.lockTreeMap.writeLock().lockInterruptibly();
            try {
            int validMsgCnt = 0;
            for (MessageExt msg : msgs) {
            MessageExt old = msgTreeMap.put(msg.getQueueOffset(), msg);
            if (null == old) {
            validMsgCnt++;
            this.queueOffsetMax = msg.getQueueOffset();
            msgSize.addAndGet(msg.getBody().length);
            }
            }
            msgCount.addAndGet(validMsgCnt);


            if (!msgTreeMap.isEmpty() && !this.consuming) {
            dispatchToConsume = true;
            this.consuming = true;
            }


            if (!msgs.isEmpty()) {
            MessageExt messageExt = msgs.get(msgs.size() - 1);
            String property = messageExt.getProperty(MessageConst.PROPERTY_MAX_OFFSET);
            if (property != null) {
            long accTotal = Long.parseLong(property) - messageExt.getQueueOffset();
            if (accTotal > 0) {
            this.msgAccCnt = accTotal;
            }
            }
            }
            } finally {
            this.lockTreeMap.writeLock().unlock();
            }
            } catch (InterruptedException e) {
            log.error("putMessage exception", e);
            }


            return dispatchToConsume;
            }


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

            评论