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..3393e41e2e1 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 a01c356f779..bfe692bf964 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 @@ -454,6 +454,18 @@ public void testClearGroup() { Mockito.eq(GROUP), Mockito.eq(TOPIC), Mockito.eq(ConfigurationManager.getProxyConfig().getInvisibleTimeMillisWhenClear())); } + @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); @@ -463,4 +475,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 +}