From 4d8997cbf8456793816eb211d0b5db898beecaea Mon Sep 17 00:00:00 2001 From: Rui <1685901819@qq.com> Date: Sat, 1 Aug 2026 18:35:27 +0800 Subject: [PATCH] [ISSUE #10739] Complete futures when proxy executors reject tasks Signed-off-by: Rui <1685901819@qq.com> --- .../rocketmq/common/utils/FutureUtils.java | 7 ++- .../common/utils/FutureUtilsTest.java | 45 +++++++++++++++++++ .../processor/DefaultMessagingProcessor.java | 6 ++- .../DefaultMessagingProcessorTest.java | 22 +++++++++ 4 files changed, 77 insertions(+), 3 deletions(-) create mode 100644 common/src/test/java/org/apache/rocketmq/common/utils/FutureUtilsTest.java diff --git a/common/src/main/java/org/apache/rocketmq/common/utils/FutureUtils.java b/common/src/main/java/org/apache/rocketmq/common/utils/FutureUtils.java index fb88b0a391f..cb18bc4abde 100644 --- a/common/src/main/java/org/apache/rocketmq/common/utils/FutureUtils.java +++ b/common/src/main/java/org/apache/rocketmq/common/utils/FutureUtils.java @@ -24,13 +24,18 @@ public class FutureUtils { public static CompletableFuture appendNextFuture(CompletableFuture future, CompletableFuture nextFuture, ExecutorService executor) { - future.whenCompleteAsync((t, throwable) -> { + CompletableFuture completionFuture = future.whenCompleteAsync((t, throwable) -> { if (throwable != null) { nextFuture.completeExceptionally(throwable); } else { nextFuture.complete(t); } }, executor); + completionFuture.whenComplete((ignored, throwable) -> { + if (throwable != null) { + nextFuture.completeExceptionally(throwable); + } + }); return nextFuture; } diff --git a/common/src/test/java/org/apache/rocketmq/common/utils/FutureUtilsTest.java b/common/src/test/java/org/apache/rocketmq/common/utils/FutureUtilsTest.java new file mode 100644 index 00000000000..a4104c95a2a --- /dev/null +++ b/common/src/test/java/org/apache/rocketmq/common/utils/FutureUtilsTest.java @@ -0,0 +1,45 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.common.utils; + +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.RejectedExecutionException; +import org.junit.Test; + +import static org.junit.Assert.assertThrows; +import static org.junit.Assert.assertTrue; + +public class FutureUtilsTest { + + @Test + public void testAddExecutorCompletesExceptionallyWhenExecutorRejectsTask() { + ExecutorService executor = Executors.newSingleThreadExecutor(); + executor.shutdown(); + + CompletableFuture result = FutureUtils.addExecutor( + CompletableFuture.completedFuture("value"), executor); + + assertTrue("the returned future must not remain pending", result.isDone()); + assertTrue(result.isCompletedExceptionally()); + CompletionException exception = assertThrows(CompletionException.class, () -> result.join()); + assertTrue(exception.getCause() instanceof RejectedExecutionException); + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java index a56bc42596b..b13ed0eccce 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java @@ -86,7 +86,8 @@ protected DefaultMessagingProcessor(ServiceManager serviceManager) { 1, TimeUnit.MINUTES, "ProducerProcessorExecutor", - proxyConfig.getProducerProcessorThreadPoolQueueCapacity() + proxyConfig.getProducerProcessorThreadPoolQueueCapacity(), + new ThreadPoolExecutor.AbortPolicy() ); this.consumerProcessorExecutor = ThreadPoolMonitor.createAndMonitor( proxyConfig.getConsumerProcessorThreadPoolNums(), @@ -94,7 +95,8 @@ protected DefaultMessagingProcessor(ServiceManager serviceManager) { 1, TimeUnit.MINUTES, "ConsumerProcessorExecutor", - proxyConfig.getConsumerProcessorThreadPoolQueueCapacity() + proxyConfig.getConsumerProcessorThreadPoolQueueCapacity(), + new ThreadPoolExecutor.AbortPolicy() ); this.serviceManager = serviceManager; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessorTest.java index 9a3ea987d09..274c0f9ae97 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessorTest.java @@ -19,7 +19,9 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.atomic.AtomicInteger; +import org.apache.rocketmq.common.utils.FutureUtils; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.apache.rocketmq.remoting.protocol.RequestCode; import org.junit.Assert; @@ -107,4 +109,24 @@ public void testRequestOnewayShouldRestoreOpaqueWhenForwardThrows() { Assert.assertEquals(originalOpaque, request.getOpaque()); } + + @Test + public void testProcessorFutureCompletesWhenExecutorIsShutDown() { + defaultMessagingProcessor.producerProcessorExecutor.shutdown(); + defaultMessagingProcessor.consumerProcessorExecutor.shutdown(); + + assertRejectedExecutor(FutureUtils.addExecutor( + CompletableFuture.completedFuture("value"), + defaultMessagingProcessor.producerProcessorExecutor)); + assertRejectedExecutor(FutureUtils.addExecutor( + CompletableFuture.completedFuture("value"), + defaultMessagingProcessor.consumerProcessorExecutor)); + } + + private void assertRejectedExecutor(CompletableFuture result) { + Assert.assertTrue("the returned future must not remain pending", result.isDone()); + Assert.assertTrue(result.isCompletedExceptionally()); + CompletionException exception = Assert.assertThrows(CompletionException.class, () -> result.join()); + Assert.assertTrue(exception.getCause() instanceof RejectedExecutionException); + } }