Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -272,6 +272,11 @@ public CompletableFuture<PopResult> popMessage(ProxyContext ctx, AddressableMess
sortMap.get(key).add(messageExt.getQueueOffset());
}
Map<String, String> map = new HashMap<>(5);
List<MessageExt> 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
Expand All @@ -285,14 +290,31 @@ public CompletableFuture<PopResult> 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<Long> sortQueueOffsets = sortMap.get(key);
List<Long> 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) {
Expand All @@ -306,6 +328,19 @@ public CompletableFuture<PopResult> 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);
Comment on lines 328 to +341
if (validMessageExtList.isEmpty() && !messageExtList.isEmpty()) {
popResult.setPopStatus(PopStatus.NO_NEW_MSG);
}
}
return popResult;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -340,6 +340,82 @@ public void testPopMessageWriteAndFlush() throws Exception {
}
}

@Test
public void testPopMessageShouldSkipMessageWithMissingOffsetMetadata() throws Exception {
Comment on lines +343 to +344
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<PopResult> 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, "");
Expand Down Expand Up @@ -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());
Expand Down
Loading