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();