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

morningman pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new c2cb912352d [bugfix](arrowflight) fix some off heap memory leak on FE 
(#66437)
c2cb912352d is described below

commit c2cb912352de1fae3f88a772be0f775497291227
Author: yiguolei <[email protected]>
AuthorDate: Wed Aug 5 09:57:23 2026 +0800

    [bugfix](arrowflight) fix some off heap memory leak on FE (#66437)
    
    ### What problem does this PR solve?
    
    Issue Number: #65615
---
 .../arrowflight/FlightSqlConnectProcessor.java        |  8 ++++----
 .../service/arrowflight/results/FlightSqlChannel.java |  1 +
 .../results/FlightSqlResultCacheEntry.java            |  2 +-
 .../arrowflight/sessions/FlightSqlConnectContext.java |  3 ---
 .../arrowflight/sessions/FlightSqlConnectPoolMgr.java | 19 +++++++++++++++++--
 .../sessions/FlightSqlConnectPoolMgrTest.java         |  7 +++++++
 6 files changed, 30 insertions(+), 10 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlConnectProcessor.java
 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlConnectProcessor.java
index aa69339fd31..24e30c3942f 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlConnectProcessor.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlConnectProcessor.java
@@ -146,11 +146,11 @@ public class FlightSqlConnectProcessor extends 
ConnectProcessor implements AutoC
                 }
                 endpointLoc.setResultPublicAccessAddr(resultPublicAccessAddr);
                 if (pResult.hasSchema() && pResult.getSchema().size() > 0) {
-                    RootAllocator rootAllocator = new 
RootAllocator(Integer.MAX_VALUE);
-                    ArrowStreamReader arrowStreamReader = new 
ArrowStreamReader(
-                            new 
ByteArrayInputStream(pResult.getSchema().toByteArray()), rootAllocator);
-                    try {
+                    try (RootAllocator rootAllocator = new 
RootAllocator(Integer.MAX_VALUE);
+                            ArrowStreamReader arrowStreamReader = new 
ArrowStreamReader(
+                                    new 
ByteArrayInputStream(pResult.getSchema().toByteArray()), rootAllocator)) {
                         Schema schema;
+                        // SchemaRoot belongs to ArrowStreamReader, it will be 
released when ArrowStreamReader is closed
                         VectorSchemaRoot root = 
arrowStreamReader.getVectorSchemaRoot();
                         List<FieldVector> fieldVectors = 
root.getFieldVectors();
                         if (fieldVectors.size() != resultOutputExprs.size()) {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java
 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java
index 4f54132876b..2781994dfa1 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java
@@ -168,6 +168,7 @@ public class FlightSqlChannel {
 
     public void close() {
         reset();
+        allocator.close();
     }
 
     private static class ResultRemovalListener implements 
RemovalListener<String, FlightSqlResultCacheEntry> {
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlResultCacheEntry.java
 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlResultCacheEntry.java
index 12ce04ca8ed..6edc868b0dd 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlResultCacheEntry.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlResultCacheEntry.java
@@ -42,7 +42,7 @@ public final class FlightSqlResultCacheEntry implements 
AutoCloseable {
 
     @Override
     public void close() throws Exception {
-        vectorSchemaRoot.clear();
+        vectorSchemaRoot.close();
     }
 
     @Override
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectContext.java
 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectContext.java
index 75f1c0ee4bf..ceddfaa5639 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectContext.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectContext.java
@@ -62,9 +62,6 @@ public class FlightSqlConnectContext extends ConnectContext {
 
     @Override
     protected void closeChannel() {
-        if (flightSqlChannel != null) {
-            flightSqlChannel.close();
-        }
         
connectScheduler.getFlightSqlConnectPoolMgr().unregisterConnection(this);
     }
 
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectPoolMgr.java
 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectPoolMgr.java
index 743f5c49d44..04982a7fedc 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectPoolMgr.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectPoolMgr.java
@@ -20,6 +20,7 @@ package org.apache.doris.service.arrowflight.sessions;
 import org.apache.doris.qe.ConnectContext;
 import org.apache.doris.qe.ConnectContext.ConnectType;
 import org.apache.doris.qe.ConnectPoolMgr;
+import org.apache.doris.service.arrowflight.results.FlightSqlChannel;
 
 import com.google.common.collect.Maps;
 import org.apache.logging.log4j.LogManager;
@@ -55,10 +56,24 @@ public class FlightSqlConnectPoolMgr extends ConnectPoolMgr 
{
 
     @Override
     public void unregisterConnection(ConnectContext ctx) {
+        // All Flight SQL session teardown paths (idle/query timeout, bearer 
token expiry, and
+        // explicit CloseSession) reach here. Release channel-cached Arrow 
results before removing
+        // the context from the pool.
+        FlightSqlChannel flightSqlChannel = ctx.getFlightSqlChannel();
+        if (flightSqlChannel != null) {
+            try {
+                flightSqlChannel.close();
+            } catch (Throwable t) {
+                // RootAllocator.close() marks the allocator closed before it 
reports outstanding
+                // bytes. The error is actionable, but session teardown must 
still release the
+                // coordinator, transaction and pool/token bookkeeping below.
+                LOG.warn("failed to close Flight SQL channel while 
unregistering connection {}, peer identity {}",
+                        ctx.getConnectionId(), ctx.getPeerIdentity(), t);
+            }
+        }
         // Finalize any Arrow Flight query whose coordinator was kept alive 
across the
         // GetFlightInfo -> DoGet phases (see #62259), releasing its resources 
(e.g. external-table
-        // batch SplitSources and the query queue slot) when the connection is 
torn down: idle/query
-        // timeout, bearer token expiry, or explicit CloseSession all reach 
here.
+        // batch SplitSources and the query queue slot).
         ctx.closeFlightSqlDeferredExecutors();
         ctx.closeTxn();
         if (connectionMap.remove(ctx.getConnectionId()) != null) {
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectPoolMgrTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectPoolMgrTest.java
index 80dcb5b4346..0513221569d 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectPoolMgrTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/sessions/FlightSqlConnectPoolMgrTest.java
@@ -18,6 +18,7 @@
 package org.apache.doris.service.arrowflight.sessions;
 
 import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.service.arrowflight.results.FlightSqlChannel;
 
 import org.junit.Assert;
 import org.junit.Test;
@@ -34,11 +35,14 @@ public class FlightSqlConnectPoolMgrTest {
     public void testUnregisterConnectionFinalizesDeferredExecutors() {
         FlightSqlConnectPoolMgr poolMgr = new FlightSqlConnectPoolMgr(100);
         ConnectContext ctx = Mockito.mock(ConnectContext.class);
+        FlightSqlChannel channel = Mockito.mock(FlightSqlChannel.class);
+        Mockito.when(ctx.getFlightSqlChannel()).thenReturn(channel);
 
         poolMgr.unregisterConnection(ctx);
 
         // The deferred coordinators must be released on teardown even though 
this connection was
         // never registered in the pool (an abandoned connection is still 
cleaned up, not leaked).
+        Mockito.verify(channel).close();
         Mockito.verify(ctx).closeFlightSqlDeferredExecutors();
     }
 
@@ -49,6 +53,8 @@ public class FlightSqlConnectPoolMgrTest {
     public void testUnregisterRegisteredConnectionFinalizesDeferredExecutors() 
{
         FlightSqlConnectPoolMgr poolMgr = new FlightSqlConnectPoolMgr(100);
         ConnectContext ctx = Mockito.mock(ConnectContext.class);
+        FlightSqlChannel channel = Mockito.mock(FlightSqlChannel.class);
+        Mockito.when(ctx.getFlightSqlChannel()).thenReturn(channel);
         Mockito.when(ctx.getConnectionId()).thenReturn(7);
         
Mockito.when(ctx.getConnectType()).thenReturn(ConnectContext.ConnectType.ARROW_FLIGHT_SQL);
         Mockito.when(ctx.getPeerIdentity()).thenReturn("token-7");
@@ -56,6 +62,7 @@ public class FlightSqlConnectPoolMgrTest {
 
         poolMgr.unregisterConnection(ctx);
 
+        Mockito.verify(channel).close();
         Mockito.verify(ctx).closeFlightSqlDeferredExecutors();
         Assert.assertNull(poolMgr.getConnectionMap().get(7));
     }


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

Reply via email to