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

yiguolei pushed a commit to branch branch-4.2
in repository https://gitbox.apache.org/repos/asf/doris.git

commit fc68393b41e88caf7714ff1582db4bad4a9156d7
Author: 924060929 <[email protected]>
AuthorDate: Wed Sep 16 14:15:56 2026 +0800

    branch-4.1: [fix](fe) Fix ReadListener leak on rejected worker task 
(#62679) (#67990)
    
    Cherry-pick of #62679 to branch-4.1.
    
    ### What problem does this PR solve?
    
    ReadListener does not handle RejectedExecutionException, which may cause
    the Client to hang when concurrency is extremely high. AcceptListener
    has already performed similar handling.
    
    ### Release note
    
    Fix ReadListener leak on rejected worker task that could cause client
    hang under high concurrency.
    
    Co-authored-by: HonestManXin <[email protected]>
---
 .../java/org/apache/doris/mysql/ReadListener.java  | 37 +++++++++++++++-------
 .../apache/doris/mysql/ConnectionExceedTest.java   | 37 ++++++++++++++++++++++
 2 files changed, 63 insertions(+), 11 deletions(-)

diff --git a/fe/fe-core/src/main/java/org/apache/doris/mysql/ReadListener.java 
b/fe/fe-core/src/main/java/org/apache/doris/mysql/ReadListener.java
index c43954a98b1..4b1c87cd60c 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/mysql/ReadListener.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/mysql/ReadListener.java
@@ -26,6 +26,8 @@ import org.xnio.ChannelListener;
 import org.xnio.XnioIoThread;
 import org.xnio.conduits.ConduitStreamSourceChannel;
 
+import java.util.concurrent.RejectedExecutionException;
+
 /**
  * listener for handle mysql cmd.
  */
@@ -46,23 +48,36 @@ public class ReadListener implements 
ChannelListener<ConduitStreamSourceChannel>
         XnioIoThread.requireCurrentThread();
         ctx.suspendAcceptQuery();
         // start async query handle in task thread.
-        channel.getWorker().execute(() -> {
-            ctx.setThreadLocalInfo();
-            try {
-                connectProcessor.processOnce();
-                if (!ctx.isKilled()) {
-                    ctx.resumeAcceptQuery();
-                } else {
-                    ctx.stopAcceptQuery();
+        try {
+            channel.getWorker().execute(() -> {
+                ctx.setThreadLocalInfo();
+                try {
+                    connectProcessor.processOnce();
+                    if (!ctx.isKilled()) {
+                        ctx.resumeAcceptQuery();
+                    } else {
+                        ctx.stopAcceptQuery();
+                        ctx.cleanup();
+                    }
+                } catch (Throwable e) {
+                    LOG.warn("Exception happened in one session(" + ctx + 
").", e);
+                    ctx.setKilled();
                     ctx.cleanup();
+                } finally {
+                    ConnectContext.remove();
                 }
-            } catch (Throwable e) {
-                LOG.warn("Exception happened in one session(" + ctx + ").", e);
+            });
+        } catch (RejectedExecutionException e) {
+            LOG.warn("Failed to submit query task for one session({}).", ctx, 
e);
+            // Keep the same ConnectContext thread-local lifecycle as the 
normal async path,
+            // so that cleanup()/close listener can access 
ConnectContext.get() if needed.
+            ctx.setThreadLocalInfo();
+            try {
                 ctx.setKilled();
                 ctx.cleanup();
             } finally {
                 ConnectContext.remove();
             }
-        });
+        }
     }
 }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/mysql/ConnectionExceedTest.java 
b/fe/fe-core/src/test/java/org/apache/doris/mysql/ConnectionExceedTest.java
index c3066ff6197..15c36b5f54d 100644
--- a/fe/fe-core/src/test/java/org/apache/doris/mysql/ConnectionExceedTest.java
+++ b/fe/fe-core/src/test/java/org/apache/doris/mysql/ConnectionExceedTest.java
@@ -22,7 +22,9 @@ import org.apache.doris.catalog.Env;
 import org.apache.doris.common.ErrorCode;
 import org.apache.doris.mysql.privilege.Auth;
 import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.ConnectProcessor;
 import org.apache.doris.qe.ConnectScheduler;
+import org.apache.doris.qe.QueryState;
 import org.apache.doris.service.ExecuteEnv;
 import 
org.apache.doris.service.arrowflight.sessions.FlightSessionsWithTokenManager;
 import org.apache.doris.service.arrowflight.tokens.FlightTokenDetails;
@@ -32,7 +34,15 @@ import mockit.Expectations;
 import mockit.Mocked;
 import org.junit.Assert;
 import org.junit.Test;
+import org.mockito.InOrder;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
 import org.xnio.StreamConnection;
+import org.xnio.XnioIoThread;
+import org.xnio.XnioWorker;
+import org.xnio.conduits.ConduitStreamSourceChannel;
+
+import java.util.concurrent.RejectedExecutionException;
 
 public class ConnectionExceedTest {
     @Mocked
@@ -105,6 +115,33 @@ public class ConnectionExceedTest {
         Assert.assertEquals(ErrorCode.ERR_TOO_MANY_USER_CONNECTIONS, 
context3.getState().getErrorCode());
     }
 
+    @Test
+    public void testHandleReadEventRejectedExecution() throws Exception {
+        try (MockedStatic<XnioIoThread> mockedIoThread = 
Mockito.mockStatic(XnioIoThread.class)) {
+            ConnectContext context = Mockito.mock(ConnectContext.class);
+            QueryState queryState = Mockito.mock(QueryState.class);
+            ConnectProcessor processor = Mockito.mock(ConnectProcessor.class);
+            ConduitStreamSourceChannel channel = 
Mockito.mock(ConduitStreamSourceChannel.class);
+            XnioWorker worker = Mockito.mock(XnioWorker.class);
+
+            Mockito.when(context.getState()).thenReturn(queryState);
+            Mockito.when(channel.getWorker()).thenReturn(worker);
+            Mockito.doThrow(new RejectedExecutionException("queue full"))
+                    .when(worker).execute(Mockito.any(Runnable.class));
+
+            ReadListener listener = new ReadListener(context, processor);
+            listener.handleEvent(channel);
+
+            InOrder contextInOrder = Mockito.inOrder(context);
+            contextInOrder.verify(context).suspendAcceptQuery();
+            contextInOrder.verify(context).setThreadLocalInfo();
+            contextInOrder.verify(context).setKilled();
+            contextInOrder.verify(context).cleanup();
+            Mockito.verifyNoInteractions(queryState);
+            Mockito.verifyNoInteractions(processor);
+        }
+    }
+
     @Test
     public void testFlightSessionConnectionExceed() throws Exception {
         // Create a scheduler with small max connections


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to