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]
