diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java index 07f16e98965..0b20050631d 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java @@ -568,7 +568,13 @@ private void reviveMsgFromCk(PopCheckPoint popCheckPoint) { // retry msg long msgOffset = popCheckPoint.ackOffsetByIndex((byte) j); CompletableFuture> future = getBizMessage(popCheckPoint, msgOffset) - .thenApply(rst -> { + .handle((rst, throwable) -> { + if (throwable != null) { + POP_LOGGER.error("reviveQueueId={}, get biz msg failed, topic:{}, qid:{}, offset:{}, brokerName:{}", + queueId, popCheckPoint.getTopic(), popCheckPoint.getQueueId(), msgOffset, + popCheckPoint.getBrokerName(), throwable); + return new Pair<>(msgOffset, false); + } MessageExt message = rst.getLeft(); if (message == null) { POP_LOGGER.info("reviveQueueId={}, can not get biz msg, topic:{}, qid:{}, offset:{}, brokerName:{}, info:{}, retry:{}, then continue", diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java index fa7e9982e1f..e109ef1b669 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/PopReviveServiceTest.java @@ -17,6 +17,7 @@ package org.apache.rocketmq.broker.processor; import com.alibaba.fastjson2.JSON; +import org.apache.commons.lang3.reflect.FieldUtils; import org.apache.commons.lang3.tuple.Triple; import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.broker.failover.EscapeBridge; @@ -38,6 +39,7 @@ import org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig; import org.apache.rocketmq.store.AppendMessageResult; import org.apache.rocketmq.store.AppendMessageStatus; +import org.apache.rocketmq.store.GetMessageResult; import org.apache.rocketmq.store.MessageStore; import org.apache.rocketmq.store.PutMessageResult; import org.apache.rocketmq.store.PutMessageStatus; @@ -58,6 +60,7 @@ import java.util.Arrays; import java.util.Collection; import java.util.List; +import java.util.NavigableMap; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; @@ -389,6 +392,66 @@ public void testReviveMsgFromCk_messageNotFound_needRetry() throws Throwable { verify(messageStore, times(1)).putMessage(any(MessageExtBrokerInner.class)); // rewrite CK } + @Test + public void testReviveMsgFromCk_getBizMessageExceptional_rewriteCK() throws Throwable { + PopCheckPoint ck = buildPopCheckPoint(0, 0, 1); + PopReviveService.ConsumeReviveObj reviveObj = new PopReviveService.ConsumeReviveObj(); + reviveObj.map.put("", ck); + reviveObj.endTime = System.currentTimeMillis(); + + ArgumentCaptor commitOffsetCaptor = ArgumentCaptor.forClass(Long.class); + doNothing().when(consumerOffsetManager).commitOffset(anyString(), anyString(), anyString(), anyInt(), + commitOffsetCaptor.capture()); + + CompletableFuture failed = new CompletableFuture<>(); + failed.completeExceptionally(new RuntimeException("store read failed")); + // Exercise the real bridge; inject the fault only at the MessageStore async boundary. + EscapeBridge realEscapeBridge = new EscapeBridge(brokerController); + when(brokerController.getEscapeBridge()).thenReturn(realEscapeBridge); + when(brokerController.getMessageStoreByBrokerName(ck.getBrokerName())).thenReturn(messageStore); + when(messageStore.getMessageAsync(anyString(), anyString(), anyInt(), anyLong(), anyInt(), any())) + .thenReturn(failed); + + popReviveService.mergeAndRevive(reviveObj); + + NavigableMap inflight = (NavigableMap) FieldUtils.readField( + popReviveService, "inflightReviveRequestMap", true); + assertEquals(1, reviveObj.newOffset); + assertEquals(1, commitOffsetCaptor.getValue().longValue()); + assertEquals(0, inflight.size()); + // An exceptional async read must retain retryability by rewriting the checkpoint. + verify(messageStore, times(1)).putMessage(any(MessageExtBrokerInner.class)); + } + + @Test + public void testReviveMsgFromCk_storeFailsAfterOffsetCommit_rewriteCK() throws Throwable { + PopCheckPoint ck = buildPopCheckPoint(0, 0, 1); + PopReviveService.ConsumeReviveObj reviveObj = new PopReviveService.ConsumeReviveObj(); + reviveObj.map.put("", ck); + reviveObj.endTime = System.currentTimeMillis(); + CompletableFuture readFuture = new CompletableFuture<>(); + EscapeBridge realEscapeBridge = new EscapeBridge(brokerController); + when(brokerController.getEscapeBridge()).thenReturn(realEscapeBridge); + when(brokerController.getMessageStoreByBrokerName(ck.getBrokerName())).thenReturn(messageStore); + when(messageStore.getMessageAsync(anyString(), anyString(), anyInt(), anyLong(), anyInt(), any())) + .thenReturn(readFuture); + + popReviveService.mergeAndRevive(reviveObj); + + NavigableMap inflight = (NavigableMap) FieldUtils.readField( + popReviveService, "inflightReviveRequestMap", true); + assertEquals(1, reviveObj.newOffset); + assertEquals(1, inflight.size()); + verify(consumerOffsetManager).commitOffset(PopAckConstants.LOCAL_HOST, PopAckConstants.REVIVE_GROUP, + REVIVE_TOPIC, REVIVE_QUEUE_ID, 1); + verify(messageStore, times(0)).putMessage(any(MessageExtBrokerInner.class)); + + readFuture.completeExceptionally(new RuntimeException("delayed store read failure")); + + assertEquals(0, inflight.size()); + verify(messageStore, times(1)).putMessage(any(MessageExtBrokerInner.class)); + } + @Test public void testReviveMsgFromCk_messageNotFound_needRetry_end() throws Throwable { brokerConfig.setSkipWhenCKRePutReachMaxTimes(true);