{"record":{"id":"c0c35166e65334a6","repo":"apache/rocketmq","slug":"failed-to-get-max-offset-in-queue-c0c351","errorCode":null,"errorMessage":"Failed to get max offset in queue","messagePattern":"Failed to get max offset in queue","errorType":"exception","errorClass":"RemotingCommandException","httpStatus":null,"severity":"error","filePath":"broker/src/main/java/org/apache/rocketmq/broker/processor/PeekMessageProcessor.java","lineNumber":239,"sourceCode":"            default:\n                assert false;\n        }\n        return response;\n    }\n\n    private long peekMsgFromQueue(boolean isRetry, GetMessageResult getMessageResult,\n        PeekMessageRequestHeader requestHeader, int queueId, long restNum, int reviveQid, Channel channel,\n        long popTime) throws RemotingCommandException {\n        String topic = isRetry ?\n            KeyBuilder.buildPopRetryTopic(requestHeader.getTopic(), requestHeader.getConsumerGroup(), brokerController.getBrokerConfig().isEnableRetryTopicV2())\n            : requestHeader.getTopic();\n        GetMessageResult getMessageTmpResult;\n        long offset = getPopOffset(topic, requestHeader.getConsumerGroup(), queueId);\n        try {\n            restNum = this.brokerController.getMessageStore().getMaxOffsetInQueue(topic, queueId) - offset + restNum;\n        } catch (ConsumeQueueException e) {\n            LOG.error(\"Failed to get max offset in queue. topic={}, queue-id={}\", topic, queueId, e);\n            throw new RemotingCommandException(\"Failed to get max offset in queue\", e);\n        }\n        if (getMessageResult.getMessageMapedList().size() >= requestHeader.getMaxMsgNums()) {\n            return restNum;\n        }\n        getMessageTmpResult = this.brokerController.getMessageStore().getMessage(requestHeader.getConsumerGroup(), topic, queueId, offset,\n            requestHeader.getMaxMsgNums() - getMessageResult.getMessageMapedList().size(), null);\n        // maybe store offset is not correct.\n        if (GetMessageStatus.OFFSET_TOO_SMALL.equals(getMessageTmpResult.getStatus()) || GetMessageStatus.OFFSET_OVERFLOW_BADLY.equals(getMessageTmpResult.getStatus())) {\n            offset = getMessageTmpResult.getNextBeginOffset();\n            getMessageTmpResult = this.brokerController.getMessageStore().getMessage(requestHeader.getConsumerGroup(), topic, queueId, offset,\n                requestHeader.getMaxMsgNums() - getMessageResult.getMessageMapedList().size(), null);\n        }\n        if (getMessageTmpResult != null) {\n            if (!getMessageTmpResult.getMessageMapedList().isEmpty() && !isRetry) {\n                Attributes attributes = this.brokerController.getBrokerMetricsManager().newAttributesBuilder()\n                    .put(LABEL_TOPIC, requestHeader.getTopic())\n                    .put(LABEL_CONSUMER_GROUP, requestHeader.getConsumerGroup())\n                    .put(LABEL_IS_SYSTEM, TopicValidator.isSystemTopic(requestHeader.getTopic()) || MixAll.isSysConsumerGroup(requestHeader.getConsumerGroup()))","sourceCodeStart":221,"sourceCodeEnd":257,"githubUrl":"https://github.com/apache/rocketmq/blob/293f5885719fc4aa3619446a1900f58ccfcfdd29/broker/src/main/java/org/apache/rocketmq/broker/processor/PeekMessageProcessor.java#L221-L257","documentation":"PeekMessageProcessor (pop-peek, used by proxy lightweight clients) computes remaining messages as maxOffsetInQueue - offset + restNum; a ConsumeQueueException from getMaxOffsetInQueue is wrapped in RemotingCommandException, failing the PEEK_MESSAGE request.","triggerScenarios":"Peek request against a topic whose consume queue can't be read — corrupt/missing queue files, store loading, RocksDB read error.","commonSituations":"Proxy peeking retry/pop-retry topics right after broker restart before store load completes; degraded disk; queue files damaged after crash.","solutions":["Wait for full store initialization on the broker, then retry the peek.","Check the logged topic/queue-id in the LOG.error line to locate the failing queue.","Repair/restore the consume queue store; verify with mqAdmin topicStatus that offsets are queryable."],"exampleFix":null,"handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"catch (RemotingCommandException e) { if (e.getCause() instanceof ConsumeQueueException) { retryPeekWithBackoff(pollIntervalMs); } else throw e; }","preventionTips":["Don't peek against a broker mid-startup; wait for it to be registered/healthy.","Proxy should fail over peek requests to a healthy replica broker."],"tags":["rocketmq","broker","pop-consumer","peek","consume-queue","proxy","storage"],"backgroundTag":null,"analyzedSha":"293f5885719fc4aa3619446a1900f58ccfcfdd29","analyzedAt":"2026-08-14T11:50:13.822Z","schemaVersion":2},"datasetVersion":"2026-08-15T22:17:37.221Z"}