From e60e409ba4b5ab8615cf85c37de18e34a2dfd0cf Mon Sep 17 00:00:00 2001 From: Rui <1685901819@qq.com> Date: Sat, 29 Aug 2026 14:26:50 +0800 Subject: [PATCH 1/2] [ISSUE #10985] Handle exceptional POP revive reads Signed-off-by: Rui <1685901819@qq.com> --- .../broker/processor/PopReviveService.java | 8 ++++- .../processor/PopReviveServiceTest.java | 29 +++++++++++++++++++ 2 files changed, 36 insertions(+), 1 deletion(-) 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..cbd702f9e8a 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; @@ -58,6 +59,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 +391,33 @@ 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")); + when(escapeBridge.getMessageAsync(anyString(), anyLong(), anyInt(), anyString(), anyBoolean())) + .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_messageNotFound_needRetry_end() throws Throwable { brokerConfig.setSkipWhenCKRePutReachMaxTimes(true); From a191a2953750e946c2ffc6d8ba1d9cb547e85036 Mon Sep 17 00:00:00 2001 From: Rui <1685901819@qq.com> Date: Sun, 13 Sep 2026 13:26:42 +0800 Subject: [PATCH 2/2] [ISSUE #10985] Exercise exceptional revive reads through EscapeBridge Inject immediate and delayed failures at the MessageStore future boundary, verifying checkpoint rewrite and inflight cleanup even after the revive offset has already been committed. Signed-off-by: Rui <1685901819@qq.com> --- .../processor/PopReviveServiceTest.java | 38 ++++++++++++++++++- 1 file changed, 36 insertions(+), 2 deletions(-) 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 cbd702f9e8a..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 @@ -39,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; @@ -402,9 +403,13 @@ public void testReviveMsgFromCk_getBizMessageExceptional_rewriteCK() throws Thro doNothing().when(consumerOffsetManager).commitOffset(anyString(), anyString(), anyString(), anyInt(), commitOffsetCaptor.capture()); - CompletableFuture> failed = new CompletableFuture<>(); + CompletableFuture failed = new CompletableFuture<>(); failed.completeExceptionally(new RuntimeException("store read failed")); - when(escapeBridge.getMessageAsync(anyString(), anyLong(), anyInt(), anyString(), anyBoolean())) + // 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); @@ -418,6 +423,35 @@ public void testReviveMsgFromCk_getBizMessageExceptional_rewriteCK() throws Thro 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);