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..8d3094b1808 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,7 @@ public CompletableFuture popMessage(ProxyContext ctx, AddressableMess sortMap.get(key).add(messageExt.getQueueOffset()); } Map map = new HashMap<>(5); + List validMessageExtList = new ArrayList<>(messageExtList.size()); for (MessageExt messageExt : messageExtList) { if (startOffsetInfo == null) { // we should set the check point info to extraInfo field , if the command is popMsg @@ -285,14 +286,26 @@ 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 startOffset = startOffsetInfo.get(key); + if (sortQueueOffsets == null || msgQueueOffsets == null || startOffset == null) { + log.warn("Pop response offset metadata is missing, key:{}", key); + continue; + } + int index = sortQueueOffsets.indexOf(messageExt.getQueueOffset()); + if (index < 0 || index >= msgQueueOffsets.size()) { + log.warn("Pop response offset metadata index is invalid, key:{}, index:{}, msgOffsetCount:{}", + key, index, msgQueueOffsets.size()); + 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(startOffset, responseHeader.getPopTime(), responseHeader.getInvisibleTime(), responseHeader.getReviveQid(), messageExt.getTopic(), messageQueue.getBrokerName(), messageExt.getQueueId(), msgQueueOffset) ); if (requestHeader.isOrder() && orderCountInfo != null) { @@ -306,7 +319,9 @@ 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); } + popResult.setMsgFoundList(validMessageExtList); } 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..f70331e1dc7 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,47 @@ 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); + 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(startOffsetStringBuilder.toString()); + responseHeader.setInvisibleTime(requestHeader.getInvisibleTime()); + responseHeader.setPopTime(popTime); + responseHeader.setReviveQid(reviveQueueId); + simpleChannelHandlerContext.writeAndFlush(response); + return null; + }); + + 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.FOUND); + assertThat(popResult.getMsgFoundList()).isEmpty(); + } + @Test public void testPopMessagePollingTimeout() throws Exception { RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.POLLING_TIMEOUT, "");