From caad3ee12792825237276291ca04eec74d4cb7d5 Mon Sep 17 00:00:00 2001 From: liuhy Date: Fri, 3 Jul 2026 11:32:49 +0800 Subject: [PATCH 1/2] Bound remoting callback executor queues to avoid OOM The remoting client and server public callback executors used fixed thread pools backed by unbounded LinkedBlockingQueue instances. Under callback backlogs this can retain tasks until the JVM runs out of heap, so both sides now use bounded queues with configurable capacities and the existing ThreadUtils executor factory. Constraint: Issue #9985 reports unbounded NettyRemotingServer publicExecutor queue growth as an OOM risk Rejected: Only fix server side | NettyRemotingClient had the same fixed-thread-pool pattern and callback role Confidence: high Scope-risk: narrow Directive: Keep remoting callback executor queues bounded unless a caller supplies an explicit custom executor Tested: mvn -pl remoting -Dtest=NettyRemotingServerTest,NettyRemotingClientTest test -Dspotbugs.skip=true -Dcheckstyle.skip=true Tested: git diff --check Related: https://github.com/apache/rocketmq/issues/9985 --- .../remoting/netty/NettyClientConfig.java | 9 ++++++ .../remoting/netty/NettyRemotingClient.java | 9 ++++-- .../remoting/netty/NettyRemotingServer.java | 9 ++++-- .../remoting/netty/NettyServerConfig.java | 9 ++++++ .../netty/NettyRemotingClientTest.java | 26 +++++++++++++++++ .../netty/NettyRemotingServerTest.java | 28 ++++++++++++++++++- 6 files changed, 85 insertions(+), 5 deletions(-) diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyClientConfig.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyClientConfig.java index 82601636403..fbba2c80543 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyClientConfig.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyClientConfig.java @@ -26,6 +26,7 @@ public class NettyClientConfig { */ private int clientWorkerThreads = NettySystemConfig.clientWorkerSize; private int clientCallbackExecutorThreads = Runtime.getRuntime().availableProcessors(); + private int clientCallbackExecutorQueueCapacity = 10000; private int clientOnewaySemaphoreValue = NettySystemConfig.CLIENT_ONEWAY_SEMAPHORE_VALUE; private int clientAsyncSemaphoreValue = NettySystemConfig.CLIENT_ASYNC_SEMAPHORE_VALUE; private int connectTimeoutMillis = NettySystemConfig.connectTimeoutMillis; @@ -99,6 +100,14 @@ public void setClientCallbackExecutorThreads(int clientCallbackExecutorThreads) this.clientCallbackExecutorThreads = clientCallbackExecutorThreads; } + public int getClientCallbackExecutorQueueCapacity() { + return clientCallbackExecutorQueueCapacity; + } + + public void setClientCallbackExecutorQueueCapacity(int clientCallbackExecutorQueueCapacity) { + this.clientCallbackExecutorQueueCapacity = clientCallbackExecutorQueueCapacity; + } + public long getChannelNotActiveInterval() { return channelNotActiveInterval; } diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java index 94d5ff9f3f6..9211c6d5f28 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java @@ -84,7 +84,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; +import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -153,8 +153,13 @@ public NettyRemotingClient(final NettyClientConfig nettyClientConfig, if (publicThreadNums <= 0) { publicThreadNums = 4; } + int queueCapacity = nettyClientConfig.getClientCallbackExecutorQueueCapacity(); + if (queueCapacity <= 0) { + queueCapacity = 10000; + } - this.publicExecutor = Executors.newFixedThreadPool(publicThreadNums, new ThreadFactoryImpl("NettyClientPublicExecutor_")); + this.publicExecutor = ThreadUtils.newThreadPoolExecutor(publicThreadNums, publicThreadNums, 0L, TimeUnit.MILLISECONDS, + new LinkedBlockingQueue<>(queueCapacity), new ThreadFactoryImpl("NettyClientPublicExecutor_")); this.scanExecutor = ThreadUtils.newThreadPoolExecutor(4, 10, 60, TimeUnit.SECONDS, new ArrayBlockingQueue<>(32), new ThreadFactoryImpl("NettyClientScan_thread_")); diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingServer.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingServer.java index ab2e208f484..0bfe02184ef 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingServer.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingServer.java @@ -86,7 +86,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; +import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; @@ -169,8 +169,13 @@ private ExecutorService buildPublicExecutor(NettyServerConfig nettyServerConfig) if (publicThreadNums <= 0) { publicThreadNums = 4; } + int queueCapacity = nettyServerConfig.getServerCallbackExecutorQueueCapacity(); + if (queueCapacity <= 0) { + queueCapacity = 10000; + } - return Executors.newFixedThreadPool(publicThreadNums, new ThreadFactoryImpl("NettyServerPublicExecutor_")); + return ThreadUtils.newThreadPoolExecutor(publicThreadNums, publicThreadNums, 0L, TimeUnit.MILLISECONDS, + new LinkedBlockingQueue<>(queueCapacity), new ThreadFactoryImpl("NettyServerPublicExecutor_")); } private ScheduledExecutorService buildScheduleExecutor() { diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyServerConfig.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyServerConfig.java index 664dee8371c..70bc99861e7 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyServerConfig.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyServerConfig.java @@ -26,6 +26,7 @@ public class NettyServerConfig implements Cloneable { private int listenPort = 0; private int serverWorkerThreads = 8; private int serverCallbackExecutorThreads = 0; + private int serverCallbackExecutorQueueCapacity = 10000; private int serverSelectorThreads = 3; private int serverOnewaySemaphoreValue = 256; private int serverAsyncSemaphoreValue = 64; @@ -99,6 +100,14 @@ public void setServerCallbackExecutorThreads(int serverCallbackExecutorThreads) this.serverCallbackExecutorThreads = serverCallbackExecutorThreads; } + public int getServerCallbackExecutorQueueCapacity() { + return serverCallbackExecutorQueueCapacity; + } + + public void setServerCallbackExecutorQueueCapacity(int serverCallbackExecutorQueueCapacity) { + this.serverCallbackExecutorQueueCapacity = serverCallbackExecutorQueueCapacity; + } + public int getServerAsyncSemaphoreValue() { return serverAsyncSemaphoreValue; } diff --git a/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingClientTest.java b/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingClientTest.java index 456e7ecdd59..57379dbb898 100644 --- a/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingClientTest.java +++ b/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingClientTest.java @@ -27,6 +27,7 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Semaphore; +import java.util.concurrent.ThreadPoolExecutor; import org.apache.rocketmq.remoting.InvokeCallback; import org.apache.rocketmq.remoting.RPCHook; import org.apache.rocketmq.remoting.common.SemaphoreReleaseOnlyOnce; @@ -37,6 +38,7 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.apache.rocketmq.remoting.protocol.RequestCode; import org.apache.rocketmq.remoting.protocol.ResponseCode; +import org.junit.After; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mock; @@ -64,11 +66,35 @@ public class NettyRemotingClientTest { @Mock private RPCHook rpcHookMock; + @After + public void tearDown() { + remotingClient.shutdown(); + } + @Test public void testSetCallbackExecutor() { ExecutorService customized = Executors.newCachedThreadPool(); remotingClient.setCallbackExecutor(customized); assertThat(remotingClient.getCallbackExecutor()).isEqualTo(customized); + customized.shutdown(); + } + + @Test + public void publicExecutorShouldUseBoundedQueue() { + NettyClientConfig nettyClientConfig = new NettyClientConfig(); + nettyClientConfig.setClientCallbackExecutorThreads(1); + nettyClientConfig.setClientCallbackExecutorQueueCapacity(3); + NettyRemotingClient client = new NettyRemotingClient(nettyClientConfig); + + try { + ThreadPoolExecutor publicExecutor = (ThreadPoolExecutor) client.getCallbackExecutor(); + + assertThat(publicExecutor.getCorePoolSize()).isEqualTo(1); + assertThat(publicExecutor.getMaximumPoolSize()).isEqualTo(1); + assertThat(publicExecutor.getQueue().remainingCapacity()).isEqualTo(3); + } finally { + client.shutdown(); + } } @Test diff --git a/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingServerTest.java b/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingServerTest.java index c69fcebd453..2ecfecc8971 100644 --- a/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingServerTest.java +++ b/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingServerTest.java @@ -23,6 +23,9 @@ import io.netty.util.Attribute; import io.netty.util.AttributeKey; import java.nio.charset.StandardCharsets; +import java.util.concurrent.ThreadPoolExecutor; +import org.junit.After; +import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -50,6 +53,11 @@ public void setUp() throws Exception { nettyRemotingServer = new NettyRemotingServer(nettyServerConfig); } + @After + public void tearDown() { + nettyRemotingServer.shutdown(); + } + @Test public void handleHAProxyTLV() { when(channel.attr(any(AttributeKey.class))).thenReturn(attribute); @@ -60,4 +68,22 @@ public void handleHAProxyTLV() { HAProxyTLV haProxyTLV = new HAProxyTLV((byte) 0xE1, content); nettyRemotingServer.handleHAProxyTLV(haProxyTLV, channel); } -} \ No newline at end of file + + @Test + public void publicExecutorShouldUseBoundedQueue() { + NettyServerConfig nettyServerConfig = new NettyServerConfig(); + nettyServerConfig.setServerCallbackExecutorThreads(1); + nettyServerConfig.setServerCallbackExecutorQueueCapacity(3); + NettyRemotingServer remotingServer = new NettyRemotingServer(nettyServerConfig); + + try { + ThreadPoolExecutor publicExecutor = (ThreadPoolExecutor) remotingServer.getCallbackExecutor(); + + Assert.assertEquals(1, publicExecutor.getCorePoolSize()); + Assert.assertEquals(1, publicExecutor.getMaximumPoolSize()); + Assert.assertEquals(3, publicExecutor.getQueue().remainingCapacity()); + } finally { + remotingServer.shutdown(); + } + } +} From fcba7ed27749e82599255dbaf7046b3c3f5c8f4c Mon Sep 17 00:00:00 2001 From: liuhy Date: Sun, 12 Jul 2026 07:23:01 -0700 Subject: [PATCH 2/2] [ISSUE #9985] Use CallerRunsPolicy for bounded remoting callback executors The bounded callback executors added for #9985 relied on the default AbortPolicy, so a full queue threw RejectedExecutionException per callback. NettyRemotingAbstract.executeInvokeCallback catches that, logs a WARN-with-stacktrace per rejection, then runs the callback inline on the caller (Netty I/O) thread anyway. Switching to an explicit CallerRunsPolicy keeps the same "run inline under backpressure" outcome without the per-rejection exception and log spam, and satisfies issue #9985's ask for a proper RejectedExecutionHandler. CallerRunsPolicy is already the convention used by the transaction message check listener. Added tests asserting the policy is CallerRunsPolicy and covering the queueCapacity<=0 fallback (previously uncovered per Codecov). Co-Authored-By: Claude --- .../remoting/netty/NettyRemotingClient.java | 4 +++- .../remoting/netty/NettyRemotingServer.java | 3 ++- .../netty/NettyRemotingClientTest.java | 18 ++++++++++++++++++ .../netty/NettyRemotingServerTest.java | 17 +++++++++++++++++ 4 files changed, 40 insertions(+), 2 deletions(-) diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java index 9211c6d5f28..6a09d281e8b 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java @@ -86,6 +86,7 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; @@ -159,7 +160,8 @@ public NettyRemotingClient(final NettyClientConfig nettyClientConfig, } this.publicExecutor = ThreadUtils.newThreadPoolExecutor(publicThreadNums, publicThreadNums, 0L, TimeUnit.MILLISECONDS, - new LinkedBlockingQueue<>(queueCapacity), new ThreadFactoryImpl("NettyClientPublicExecutor_")); + new LinkedBlockingQueue<>(queueCapacity), new ThreadFactoryImpl("NettyClientPublicExecutor_"), + new ThreadPoolExecutor.CallerRunsPolicy()); this.scanExecutor = ThreadUtils.newThreadPoolExecutor(4, 10, 60, TimeUnit.SECONDS, new ArrayBlockingQueue<>(32), new ThreadFactoryImpl("NettyClientScan_thread_")); diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingServer.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingServer.java index 0bfe02184ef..bef8fb16aa5 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingServer.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingServer.java @@ -175,7 +175,8 @@ private ExecutorService buildPublicExecutor(NettyServerConfig nettyServerConfig) } return ThreadUtils.newThreadPoolExecutor(publicThreadNums, publicThreadNums, 0L, TimeUnit.MILLISECONDS, - new LinkedBlockingQueue<>(queueCapacity), new ThreadFactoryImpl("NettyServerPublicExecutor_")); + new LinkedBlockingQueue<>(queueCapacity), new ThreadFactoryImpl("NettyServerPublicExecutor_"), + new ThreadPoolExecutor.CallerRunsPolicy()); } private ScheduledExecutorService buildScheduleExecutor() { diff --git a/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingClientTest.java b/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingClientTest.java index 57379dbb898..06a7bedc819 100644 --- a/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingClientTest.java +++ b/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingClientTest.java @@ -92,6 +92,24 @@ public void publicExecutorShouldUseBoundedQueue() { assertThat(publicExecutor.getCorePoolSize()).isEqualTo(1); assertThat(publicExecutor.getMaximumPoolSize()).isEqualTo(1); assertThat(publicExecutor.getQueue().remainingCapacity()).isEqualTo(3); + assertThat(publicExecutor.getRejectedExecutionHandler()) + .isInstanceOf(ThreadPoolExecutor.CallerRunsPolicy.class); + } finally { + client.shutdown(); + } + } + + @Test + public void publicExecutorShouldFallBackToDefaultQueueCapacityWhenMisconfigured() { + NettyClientConfig nettyClientConfig = new NettyClientConfig(); + nettyClientConfig.setClientCallbackExecutorThreads(1); + nettyClientConfig.setClientCallbackExecutorQueueCapacity(0); + NettyRemotingClient client = new NettyRemotingClient(nettyClientConfig); + + try { + ThreadPoolExecutor publicExecutor = (ThreadPoolExecutor) client.getCallbackExecutor(); + + assertThat(publicExecutor.getQueue().remainingCapacity()).isEqualTo(10000); } finally { client.shutdown(); } diff --git a/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingServerTest.java b/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingServerTest.java index 2ecfecc8971..d6e26d480b8 100644 --- a/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingServerTest.java +++ b/remoting/src/test/java/org/apache/rocketmq/remoting/netty/NettyRemotingServerTest.java @@ -82,6 +82,23 @@ public void publicExecutorShouldUseBoundedQueue() { Assert.assertEquals(1, publicExecutor.getCorePoolSize()); Assert.assertEquals(1, publicExecutor.getMaximumPoolSize()); Assert.assertEquals(3, publicExecutor.getQueue().remainingCapacity()); + Assert.assertTrue(publicExecutor.getRejectedExecutionHandler() instanceof ThreadPoolExecutor.CallerRunsPolicy); + } finally { + remotingServer.shutdown(); + } + } + + @Test + public void publicExecutorShouldFallBackToDefaultQueueCapacityWhenMisconfigured() { + NettyServerConfig nettyServerConfig = new NettyServerConfig(); + nettyServerConfig.setServerCallbackExecutorThreads(1); + nettyServerConfig.setServerCallbackExecutorQueueCapacity(0); + NettyRemotingServer remotingServer = new NettyRemotingServer(nettyServerConfig); + + try { + ThreadPoolExecutor publicExecutor = (ThreadPoolExecutor) remotingServer.getCallbackExecutor(); + + Assert.assertEquals(10000, publicExecutor.getQueue().remainingCapacity()); } finally { remotingServer.shutdown(); }