This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new d34adde0306 branch-4.1: [fix](fe) Fix ReadListener leak on rejected
worker task (#62679) (#67990)
d34adde0306 is described below
commit d34adde030635f4a0c0d27168a7dea6be1cb17cf
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]