From 0fbcec5f6f7deba750d5da9d2d7042f0e0422367 Mon Sep 17 00:00:00 2001 From: liuhy Date: Sat, 15 Aug 2026 01:26:40 -0700 Subject: [PATCH 1/2] fix(proxy): retain receipt handles after failed offline cleanup Signed-off-by: liuhy --- .../receipt/DefaultReceiptHandleManager.java | 2 +- .../DefaultReceiptHandleManagerTest.java | 22 ++++++++++++++++++- 2 files changed, 22 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManager.java index f9dfd825337..8692c25bc37 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManager.java @@ -266,7 +266,7 @@ private void returnHandleGroup(ReceiptHandleGroupKey key, ReceiptHandleGroup han handleGroup.computeIfPresent(msgID, handle, messageReceiptHandle -> { CompletableFuture future = new CompletableFuture<>(); eventListener.fireEvent(new RenewEvent(key, messageReceiptHandle, proxyConfig.getInvisibleTimeMillisWhenClear(), RenewEvent.EventType.CLEAR_GROUP, future)); - return CompletableFuture.completedFuture(null); + return future.thenApply(ackResult -> null); }, 0); } catch (Exception e) { log.error("error when clear handle for group. key:{}", key, e); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManagerTest.java index a01c356f779..34cc984aaef 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManagerTest.java @@ -60,6 +60,7 @@ import static org.awaitility.Awaitility.await; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; public class DefaultReceiptHandleManagerTest extends BaseServiceTest { @@ -454,6 +455,25 @@ public void testClearGroup() { Mockito.eq(GROUP), Mockito.eq(TOPIC), Mockito.eq(ConfigurationManager.getProxyConfig().getInvisibleTimeMillisWhenClear())); } + @Test + public void testClearGroupRetainsHandleWhenRenewalFails() { + Channel channel = PROXY_CONTEXT.getVal(ContextVariable.CHANNEL); + ReceiptHandleGroupKey key = new ReceiptHandleGroupKey(channel, GROUP); + receiptHandleManager.addReceiptHandle(PROXY_CONTEXT, channel, GROUP, MSG_ID, messageReceiptHandle); + CompletableFuture failedFuture = new CompletableFuture<>(); + failedFuture.completeExceptionally(new MQClientException(0, "renew failed")); + Mockito.when(messagingProcessor.changeInvisibleTime(Mockito.any(ProxyContext.class), Mockito.any(ReceiptHandle.class), + Mockito.eq(MESSAGE_ID), Mockito.eq(GROUP), Mockito.eq(TOPIC), + Mockito.eq(ConfigurationManager.getProxyConfig().getInvisibleTimeMillisWhenClear()))) + .thenReturn(failedFuture); + + receiptHandleManager.clearGroup(key); + + await().atMost(Duration.ofSeconds(1)).until(() -> + receiptHandleManager.receiptHandleGroupMap.containsKey(key)); + assertFalse(receiptHandleManager.receiptHandleGroupMap.get(key).isEmpty()); + } + @Test public void testClientOffline() { ArgumentCaptor listenerArgumentCaptor = ArgumentCaptor.forClass(ConsumerIdsChangeListener.class); @@ -463,4 +483,4 @@ public void testClientOffline() { listenerArgumentCaptor.getValue().handle(ConsumerGroupEvent.CLIENT_UNREGISTER, GROUP, new ClientChannelInfo(channel, "", LanguageCode.JAVA, 0)); assertTrue(receiptHandleManager.receiptHandleGroupMap.isEmpty()); } -} \ No newline at end of file +} From 5c34f63dd1b95d55669879e9cd661a06c10de934 Mon Sep 17 00:00:00 2001 From: liuhy Date: Sat, 15 Aug 2026 08:07:04 -0700 Subject: [PATCH 2/2] fix(proxy): retain rejected receipt cleanup tasks Signed-off-by: liuhy --- .../receipt/DefaultReceiptHandleManager.java | 13 +++++++++++-- .../receipt/DefaultReceiptHandleManagerTest.java | 12 ++++++++++++ 2 files changed, 23 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManager.java index 8692c25bc37..781d6340978 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManager.java @@ -24,6 +24,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; @@ -89,7 +90,8 @@ public DefaultReceiptHandleManager(MetadataService metadataService, ConsumerMana proxyConfig.getReturnHandleGroupThreadPoolNums() * 2, 1, TimeUnit.MINUTES, "ReturnHandleGroupWorkerThread", - proxyConfig.getRenewThreadPoolQueueCapacity() + proxyConfig.getRenewThreadPoolQueueCapacity(), + new ThreadPoolExecutor.AbortPolicy() ); consumerManager.appendConsumerIdsChangeListener(new ConsumerIdsChangeListener() { @Override @@ -251,7 +253,14 @@ protected void clearGroup(ReceiptHandleGroupKey key) { return; } ReceiptHandleGroup handleGroup = receiptHandleGroupMap.remove(key); - returnHandleGroupWorkerService.submit(() -> returnHandleGroup(key, handleGroup)); + try { + returnHandleGroupWorkerService.submit(() -> returnHandleGroup(key, handleGroup)); + } catch (RejectedExecutionException e) { + if (handleGroup != null) { + receiptHandleGroupMap.putIfAbsent(key, handleGroup); + } + log.warn("submit clear handle group task failed, will retry in the next schedule. key:{}", key, e); + } } // There is no longer any waiting for lock, and only the locked handles will be processed immediately, diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManagerTest.java index 34cc984aaef..abe8bc2099b 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/receipt/DefaultReceiptHandleManagerTest.java @@ -474,6 +474,18 @@ public void testClearGroupRetainsHandleWhenRenewalFails() { assertFalse(receiptHandleManager.receiptHandleGroupMap.get(key).isEmpty()); } + @Test + public void testClearGroupRetainsHandlesWhenCleanupSubmissionIsRejected() { + Channel channel = PROXY_CONTEXT.getVal(ContextVariable.CHANNEL); + ReceiptHandleGroupKey key = new ReceiptHandleGroupKey(channel, GROUP); + receiptHandleManager.addReceiptHandle(PROXY_CONTEXT, channel, GROUP, MSG_ID, messageReceiptHandle); + receiptHandleManager.returnHandleGroupWorkerService.shutdownNow(); + + receiptHandleManager.clearGroup(key); + + assertTrue(receiptHandleManager.receiptHandleGroupMap.containsKey(key)); + } + @Test public void testClientOffline() { ArgumentCaptor listenerArgumentCaptor = ArgumentCaptor.forClass(ConsumerIdsChangeListener.class);