This is an automated email from the ASF dual-hosted git repository.

lizhimins pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git


The following commit(s) were added to refs/heads/develop by this push:
     new df5b53d0e5 [ISSUE #11080] Skip unused remoting client event executor 
(#11081)
df5b53d0e5 is described below

commit df5b53d0e5cac20db964ca9a94300f2d63758eff
Author: qianye <[email protected]>
AuthorDate: Thu Sep 10 13:57:15 2026 +0800

    [ISSUE #11080] Skip unused remoting client event executor (#11081)
---
 .../remoting/netty/NettyRemotingClient.java        |  4 ++-
 .../remoting/netty/NettyRemotingClientTest.java    | 37 ++++++++++++++++++++++
 2 files changed, 40 insertions(+), 1 deletion(-)

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 94d5ff9f3f..627d6255f5 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
@@ -242,7 +242,9 @@ public class NettyRemotingClient extends 
NettyRemotingAbstract implements Remoti
             handler.option(ChannelOption.ALLOCATOR, 
PooledByteBufAllocator.DEFAULT);
         }
 
-        nettyEventExecutor.start();
+        if (channelEventListener != null) {
+            nettyEventExecutor.start();
+        }
 
         TimerTask timerTaskScanResponseTable = new TimerTask() {
             @Override
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 456e7ecdd5..8c9bb8237f 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,8 @@ import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Semaphore;
+import org.apache.rocketmq.common.ServiceThread;
+import org.apache.rocketmq.remoting.ChannelEventListener;
 import org.apache.rocketmq.remoting.InvokeCallback;
 import org.apache.rocketmq.remoting.RPCHook;
 import org.apache.rocketmq.remoting.common.SemaphoreReleaseOnlyOnce;
@@ -54,6 +56,7 @@ import static org.mockito.Mockito.doReturn;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.timeout;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 
@@ -64,6 +67,40 @@ public class NettyRemotingClientTest {
     @Mock
     private RPCHook rpcHookMock;
 
+    @Test
+    public void testStartWithoutChannelEventListener() throws Exception {
+        Field threadField = ServiceThread.class.getDeclaredField("thread");
+        threadField.setAccessible(true);
+        try {
+            remotingClient.start();
+            
assertThat(threadField.get(remotingClient.nettyEventExecutor)).isNull();
+        } finally {
+            remotingClient.shutdown();
+        }
+        
assertThat(threadField.get(remotingClient.nettyEventExecutor)).isNull();
+    }
+
+    @Test
+    public void testStartWithChannelEventListener() throws Exception {
+        ChannelEventListener listener = mock(ChannelEventListener.class);
+        NettyRemotingClient client = new NettyRemotingClient(new 
NettyClientConfig(), listener);
+        Channel channel = mock(Channel.class);
+        Field threadField = ServiceThread.class.getDeclaredField("thread");
+        threadField.setAccessible(true);
+        Thread eventThread;
+        try {
+            client.start();
+            eventThread = (Thread) threadField.get(client.nettyEventExecutor);
+            assertThat(eventThread).isNotNull();
+            assertThat(eventThread.isAlive()).isTrue();
+            client.putNettyEvent(new NettyEvent(NettyEventType.ACTIVE, 
"127.0.0.1:10911", channel));
+            verify(listener, timeout(3000)).onChannelActive("127.0.0.1:10911", 
channel);
+        } finally {
+            client.shutdown();
+        }
+        assertThat(eventThread.isAlive()).isFalse();
+    }
+
     @Test
     public void testSetCallbackExecutor() {
         ExecutorService customized = Executors.newCachedThreadPool();

Reply via email to