diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java index c93fa93983c..2392ad9b5a0 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java @@ -272,6 +272,11 @@ public CompletableFuture popMessage(ProxyContext ctx, AddressableMess sortMap.get(key).add(messageExt.getQueueOffset()); } Map map = new HashMap<>(5); + List validMessageExtList = new ArrayList<>(messageExtList.size()); + int missingOffsetMetadataCount = 0; + String firstMissingOffsetMetadataKey = null; + int invalidOffsetIndexCount = 0; + String firstInvalidOffsetIndexKey = null; for (MessageExt messageExt : messageExtList) { if (startOffsetInfo == null) { // we should set the check point info to extraInfo field , if the command is popMsg @@ -285,14 +290,31 @@ public CompletableFuture popMessage(ProxyContext ctx, AddressableMess } else { if (messageExt.getProperty(MessageConst.PROPERTY_POP_CK) == null) { String key = ExtraInfoUtil.getStartOffsetInfoMapKey(messageExt.getTopic(), messageExt.getQueueId()); - int index = sortMap.get(key).indexOf(messageExt.getQueueOffset()); - Long msgQueueOffset = msgOffsetInfo.get(key).get(index); + List sortQueueOffsets = sortMap.get(key); + List msgQueueOffsets = msgOffsetInfo == null ? null : msgOffsetInfo.get(key); + Long startOffsetForQueue = startOffsetInfo.get(key); + if (sortQueueOffsets == null || msgQueueOffsets == null || startOffsetForQueue == null) { + missingOffsetMetadataCount++; + if (firstMissingOffsetMetadataKey == null) { + firstMissingOffsetMetadataKey = key; + } + continue; + } + int index = sortQueueOffsets.indexOf(messageExt.getQueueOffset()); + if (index < 0 || index >= msgQueueOffsets.size()) { + invalidOffsetIndexCount++; + if (firstInvalidOffsetIndexKey == null) { + firstInvalidOffsetIndexKey = key; + } + continue; + } + Long msgQueueOffset = msgQueueOffsets.get(index); if (msgQueueOffset != messageExt.getQueueOffset()) { log.warn("Queue offset [{}] of msg is strange, not equal to the stored in msg, {}", msgQueueOffset, messageExt); } messageExt.getProperties().put(MessageConst.PROPERTY_POP_CK, - ExtraInfoUtil.buildExtraInfo(startOffsetInfo.get(key), responseHeader.getPopTime(), responseHeader.getInvisibleTime(), + ExtraInfoUtil.buildExtraInfo(startOffsetForQueue, responseHeader.getPopTime(), responseHeader.getInvisibleTime(), responseHeader.getReviveQid(), messageExt.getTopic(), messageQueue.getBrokerName(), messageExt.getQueueId(), msgQueueOffset) ); if (requestHeader.isOrder() && orderCountInfo != null) { @@ -306,6 +328,19 @@ public CompletableFuture popMessage(ProxyContext ctx, AddressableMess messageExt.getProperties().computeIfAbsent(MessageConst.PROPERTY_FIRST_POP_TIME, k -> String.valueOf(responseHeader.getPopTime())); messageExt.setBrokerName(messageQueue.getBrokerName()); messageExt.setTopic(messageQueue.getTopic()); + validMessageExtList.add(messageExt); + } + if (missingOffsetMetadataCount > 0) { + log.warn("Skipped {} POP messages because offset metadata is missing, first key:{}", + missingOffsetMetadataCount, firstMissingOffsetMetadataKey); + } + if (invalidOffsetIndexCount > 0) { + log.warn("Skipped {} POP messages because offset metadata index is invalid, first key:{}", + invalidOffsetIndexCount, firstInvalidOffsetIndexKey); + } + popResult.setMsgFoundList(validMessageExtList); + if (validMessageExtList.isEmpty() && !messageExtList.isEmpty()) { + popResult.setPopStatus(PopStatus.NO_NEW_MSG); } } return popResult; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/message/LocalMessageServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/message/LocalMessageServiceTest.java index 52ba521f802..beb30af9e4d 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/message/LocalMessageServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/message/LocalMessageServiceTest.java @@ -340,6 +340,82 @@ public void testPopMessageWriteAndFlush() throws Exception { } } + @Test + public void testPopMessageShouldSkipMessageWithMissingOffsetMetadata() throws Exception { + int reviveQueueId = 1; + long popTime = System.currentTimeMillis(); + long invisibleTime = 3000L; + long startOffset = 100L; + StringBuilder startOffsetStringBuilder = new StringBuilder(); + ExtraInfoUtil.buildStartOffsetInfo(startOffsetStringBuilder, topic, queueId, startOffset); + MessageExt message = buildMessageExt(topic, queueId, startOffset); + byte[] body = MessageDecoder.encode(message, false); + PopMessageRequestHeader requestHeader = new PopMessageRequestHeader(); + requestHeader.setInvisibleTime(invisibleTime); + mockPopMessageResponse(body, startOffsetStringBuilder.toString(), null, popTime, + requestHeader.getInvisibleTime(), reviveQueueId); + + MessageQueue messageQueue = new MessageQueue(topic, brokerName, queueId); + CompletableFuture future = localMessageService.popMessage(proxyContext, + new AddressableMessageQueue(messageQueue, ""), requestHeader, 1000L); + PopResult popResult = future.get(); + + assertThat(popResult.getPopStatus()).isEqualTo(PopStatus.NO_NEW_MSG); + assertThat(popResult.getMsgFoundList()).isEmpty(); + } + + @Test + public void testPopMessageShouldSkipMessageWithMissingStartOffset() throws Exception { + int reviveQueueId = 1; + long popTime = System.currentTimeMillis(); + long invisibleTime = 3000L; + long startOffset = 100L; + MessageExt message = buildMessageExt(topic, queueId, startOffset); + StringBuilder startOffsetInfo = new StringBuilder(); + StringBuilder msgOffsetInfo = new StringBuilder(); + ExtraInfoUtil.buildStartOffsetInfo(startOffsetInfo, topic, queueId + 1, startOffset); + ExtraInfoUtil.buildMsgOffsetInfo(msgOffsetInfo, topic, queueId, Collections.singletonList(startOffset)); + PopMessageRequestHeader requestHeader = new PopMessageRequestHeader(); + requestHeader.setInvisibleTime(invisibleTime); + mockPopMessageResponse(MessageDecoder.encode(message, false), startOffsetInfo.toString(), msgOffsetInfo.toString(), + popTime, invisibleTime, reviveQueueId); + + PopResult popResult = localMessageService.popMessage(proxyContext, + new AddressableMessageQueue(new MessageQueue(topic, brokerName, queueId), ""), requestHeader, 1000L).get(); + + assertThat(popResult.getPopStatus()).isEqualTo(PopStatus.NO_NEW_MSG); + assertThat(popResult.getMsgFoundList()).isEmpty(); + } + + @Test + public void testPopMessageShouldSkipMessageWithInvalidOffsetIndex() throws Exception { + int reviveQueueId = 1; + long popTime = System.currentTimeMillis(); + long invisibleTime = 3000L; + long startOffset = 100L; + MessageExt firstMessage = buildMessageExt(topic, queueId, startOffset); + MessageExt secondMessage = buildMessageExt(topic, queueId, startOffset + 1); + byte[] firstMessageBody = MessageDecoder.encode(firstMessage, false); + byte[] secondMessageBody = MessageDecoder.encode(secondMessage, false); + ByteBuffer body = ByteBuffer.allocate(firstMessageBody.length + secondMessageBody.length); + body.put(firstMessageBody).put(secondMessageBody); + StringBuilder startOffsetInfo = new StringBuilder(); + StringBuilder msgOffsetInfo = new StringBuilder(); + ExtraInfoUtil.buildStartOffsetInfo(startOffsetInfo, topic, queueId, startOffset); + ExtraInfoUtil.buildMsgOffsetInfo(msgOffsetInfo, topic, queueId, Collections.singletonList(startOffset)); + PopMessageRequestHeader requestHeader = new PopMessageRequestHeader(); + requestHeader.setInvisibleTime(invisibleTime); + mockPopMessageResponse(body.array(), startOffsetInfo.toString(), msgOffsetInfo.toString(), popTime, + invisibleTime, reviveQueueId); + + PopResult popResult = localMessageService.popMessage(proxyContext, + new AddressableMessageQueue(new MessageQueue(topic, brokerName, queueId), ""), requestHeader, 1000L).get(); + + assertThat(popResult.getPopStatus()).isEqualTo(PopStatus.FOUND); + assertThat(popResult.getMsgFoundList()).hasSize(1); + assertThat(popResult.getMsgFoundList().get(0).getQueueOffset()).isEqualTo(startOffset); + } + @Test public void testPopMessagePollingTimeout() throws Exception { RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.POLLING_TIMEOUT, ""); @@ -477,6 +553,30 @@ private MessageExt buildMessageExt(String topic, int queueId, long queueOffset) return message1; } + private void mockPopMessageResponse(byte[] body, String startOffsetInfo, String msgOffsetInfo, long popTime, + long invisibleTime, int reviveQueueId) throws RemotingCommandException { + Mockito.when(popMessageProcessorMock.processRequest(Mockito.any(SimpleChannelHandlerContext.class), Mockito.argThat(argument -> { + boolean first = argument.getCode() == RequestCode.POP_MESSAGE; + boolean second = argument.readCustomHeader() instanceof PopMessageRequestHeader; + return first && second; + }))).thenAnswer(invocation -> { + SimpleChannelHandlerContext simpleChannelHandlerContext = invocation.getArgument(0); + RemotingCommand request = invocation.getArgument(1); + RemotingCommand response = RemotingCommand.createResponseCommand(PopMessageResponseHeader.class); + response.setOpaque(request.getOpaque()); + response.setCode(ResponseCode.SUCCESS); + response.setBody(body); + PopMessageResponseHeader responseHeader = (PopMessageResponseHeader) response.readCustomHeader(); + responseHeader.setStartOffsetInfo(startOffsetInfo); + responseHeader.setMsgOffsetInfo(msgOffsetInfo); + responseHeader.setInvisibleTime(invisibleTime); + responseHeader.setPopTime(popTime); + responseHeader.setReviveQid(reviveQueueId); + simpleChannelHandlerContext.writeAndFlush(response); + return null; + }); + } + private void assertMessageExt(MessageExt messageExt1, MessageExt messageExt2) { assertThat(messageExt1.getBody()).isEqualTo(messageExt2.getBody()); assertThat(messageExt1.getTopic()).isEqualTo(messageExt2.getTopic());