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 9a5afd4958f [fix](arrow-flight-sql) Close temporary VectorSchemaRoot 
in createPreparedStatement to fix FE direct memory leak (#65311)
9a5afd4958f is described below

commit 9a5afd4958f2794a2dd458b2e7463315e9746ce3
Author: OneG <[email protected]>
AuthorDate: Fri Jul 17 14:08:40 2026 +0800

    [fix](arrow-flight-sql) Close temporary VectorSchemaRoot in 
createPreparedStatement to fix FE direct memory leak (#65311)
    
    ```
    ### What problem does this PR solve?
    
    Issue Number: close #65305
    
    Problem Summary:
    
    `DorisFlightSqlProducer#createPreparedStatement` creates two 
`VectorSchemaRoot`
    instances (an empty one for the parameter schema, and one from
    `FlightSqlChannel#createOneOneSchemaRoot(...)` for the result metadata),
    extracts only their `Schema`, and never closes them.
    
    `createOneOneSchemaRoot` allocates a `VarCharVector` from the channel
    `RootAllocator` (off-heap, Netty pooled direct buffer). Since the returned
    `VectorSchemaRoot` is never closed, its off-heap buffer is leaked on every
    prepare call (effectively every Arrow Flight query, because ADBC prepares 
each
    statement). The buffer cannot be reclaimed by GC (strongly referenced by the
    allocator) nor by closing the client session (`FlightSqlChannel#close()` 
only
    invalidates the result cache and does not close the allocator; the leaked 
root
    is never stored in any session map).
    
    Under continuous Arrow Flight query load this makes FE direct memory grow
    monotonically until:
    `java.lang.OutOfMemoryError: Cannot reserve ... bytes of direct buffer 
memory (Internal; Prepare)`
    
    This PR wraps both temporary roots in try-with-resources so their off-heap
    buffers are released after the `Schema` is extracted. `Schema` is an 
immutable
    POJO and remains valid after the root is closed. Compare with
    `getCatalogs`/`getSchemas`/`getTables` in the same producer, which already 
use
    try-with-resources.
    
    ---------
    
    Signed-off-by: 柳吟风 <[email protected]>
    Signed-off-by: morningman <[email protected]>
    Co-authored-by: 柳吟风 <[email protected]>
    Co-authored-by: morningman <[email protected]>
---
 .../arrowflight/DorisFlightSqlProducer.java        |  16 ++-
 .../arrowflight/results/FlightSqlChannel.java      |   6 +
 .../arrowflight/DorisFlightSqlProducerTest.java    | 145 +++++++++++++++++++++
 3 files changed, 163 insertions(+), 4 deletions(-)

diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/DorisFlightSqlProducer.java
 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/DorisFlightSqlProducer.java
index ebc3347f07f..e2d197cb537 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/DorisFlightSqlProducer.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/DorisFlightSqlProducer.java
@@ -356,14 +356,22 @@ public class DorisFlightSqlProducer implements 
FlightSqlProducer, AutoCloseable
                 final ByteString handle = 
ByteString.copyFromUtf8(context.peerIdentity() + ":" + preparedStatementId);
                 connectContext.addPreparedQuery(preparedStatementId, query);
 
-                VectorSchemaRoot emptyVectorSchemaRoot = new 
VectorSchemaRoot(new ArrayList<>(), new ArrayList<>());
-                final Schema parameterSchema = 
emptyVectorSchemaRoot.getSchema();
+                // Close the temporary VectorSchemaRoot after extracting its 
Schema, otherwise the
+                // off-heap buffers backing its vectors are leaked on every 
prepare (FE direct memory leak).
+                final Schema parameterSchema;
+                try (VectorSchemaRoot emptyVectorSchemaRoot =
+                        new VectorSchemaRoot(new ArrayList<>(), new 
ArrayList<>())) {
+                    parameterSchema = emptyVectorSchemaRoot.getSchema();
+                }
                 // TODO FE does not have the ability to convert root fragment 
output expr into arrow schema.
                 // However, the metaData schema returned by 
createPreparedStatement is usually not used by the client,
                 // but it cannot be empty, otherwise it will be mistaken by 
the client as an updata statement.
                 // see: https://github.com/apache/arrow/issues/38911
-                Schema metaData = connectContext.getFlightSqlChannel()
-                        .createOneOneSchemaRoot("ResultMeta", 
"UNIMPLEMENTED").getSchema();
+                final Schema metaData;
+                try (VectorSchemaRoot metaSchemaRoot = 
connectContext.getFlightSqlChannel()
+                        .createOneOneSchemaRoot("ResultMeta", 
"UNIMPLEMENTED")) {
+                    metaData = metaSchemaRoot.getSchema();
+                }
                 listener.onNext(new Result(
                         Any.pack(buildCreatePreparedStatementResult(handle, 
parameterSchema, metaData)).toByteArray()));
             } catch (Exception e) {
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 8ce5f72b400..4f54132876b 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
@@ -156,6 +156,12 @@ public class FlightSqlChannel {
         return resultCache.size();
     }
 
+    // Returns the off-heap memory (bytes) currently held by this channel's 
Arrow allocator.
+    // Exposed so tests can assert prepared-statement handling does not leak 
VectorSchemaRoot buffers.
+    public long getAllocatedMemory() {
+        return allocator.getAllocatedMemory();
+    }
+
     public void reset() {
         resultCache.invalidateAll();
     }
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/DorisFlightSqlProducerTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/DorisFlightSqlProducerTest.java
new file mode 100644
index 00000000000..3a4b0facace
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/service/arrowflight/DorisFlightSqlProducerTest.java
@@ -0,0 +1,145 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.service.arrowflight;
+
+import org.apache.doris.common.FeConstants;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.service.arrowflight.results.FlightSqlChannel;
+import org.apache.doris.service.arrowflight.sessions.FlightSessionsManager;
+import org.apache.doris.service.arrowflight.sessions.FlightSqlConnectContext;
+
+import org.apache.arrow.flight.FlightProducer.CallContext;
+import org.apache.arrow.flight.FlightProducer.StreamListener;
+import org.apache.arrow.flight.Location;
+import org.apache.arrow.flight.Result;
+import 
org.apache.arrow.flight.sql.impl.FlightSql.ActionCreatePreparedStatementRequest;
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+public class DorisFlightSqlProducerTest {
+
+    private boolean prevRunningUnitTest;
+
+    @Before
+    public void setUp() {
+        // FlightSqlConnectContext.init() only reaches Env when this is false; 
keep it true so the
+        // context can be built without a running FE.
+        prevRunningUnitTest = FeConstants.runningUnitTest;
+        FeConstants.runningUnitTest = true;
+    }
+
+    @After
+    public void tearDown() {
+        FeConstants.runningUnitTest = prevRunningUnitTest;
+    }
+
+    /**
+     * Regression test for the FE direct-memory leak in {@code 
createPreparedStatement}
+     * (issue apache/doris#65305, fixed by PR #65311).
+     *
+     * <p>Before the fix, each prepare allocated a {@link 
org.apache.arrow.vector.VectorSchemaRoot}
+     * via {@code FlightSqlChannel.createOneOneSchemaRoot("ResultMeta", ...)} 
and read only its
+     * {@code Schema}, never closing the root, leaking one off-heap {@code 
VarCharVector} buffer per
+     * prepare. Because Arrow Flight has no parameter binding, every client 
query triggers a fresh
+     * prepare, so the leak grows monotonically until {@code 
MaxDirectMemorySize} is exhausted.
+     *
+     * <p>This test drives {@code createPreparedStatement} against a real 
{@link FlightSqlChannel} many
+     * times and asserts the channel's Arrow allocator holds zero bytes 
afterwards. It goes red if the
+     * try-with-resources that closes the temporary roots is removed.
+     */
+    @Test
+    public void createPreparedStatementDoesNotLeakChannelAllocator() throws 
Exception {
+        // A real flight session context owns a real FlightSqlChannel (and 
thus a real Arrow allocator),
+        // so allocator bookkeeping is exercised for real instead of mocked 
away.
+        FlightSqlConnectContext connectContext = new 
FlightSqlConnectContext("test-peer-identity");
+        FlightSqlChannel channel = connectContext.getFlightSqlChannel();
+        Assert.assertEquals("channel allocator should start empty", 0L, 
channel.getAllocatedMemory());
+
+        FlightSessionsManager sessionsManager = new FlightSessionsManager() {
+            @Override
+            public ConnectContext getConnectContext(String peerIdentity) {
+                return connectContext;
+            }
+
+            @Override
+            public ConnectContext createConnectContext(String peerIdentity) {
+                return connectContext;
+            }
+
+            @Override
+            public void closeConnectContext(String peerIdentity) {
+                // not exercised by this test
+            }
+        };
+        DorisFlightSqlProducer producer =
+                new 
DorisFlightSqlProducer(Location.forGrpcInsecure("127.0.0.1", 9090), 
sessionsManager);
+
+        CallContext callContext = Mockito.mock(CallContext.class);
+        
Mockito.when(callContext.peerIdentity()).thenReturn("test-peer-identity");
+
+        final int rounds = 100;
+        AtomicInteger errors = new AtomicInteger(0);
+        try {
+            // createPreparedStatement runs asynchronously and mutates a 
shared, non-thread-safe
+            // ConnectContext, so drive it serially: each prepare completes 
before the next is issued.
+            // The leak, if any, still accumulates on the single channel 
allocator across rounds.
+            for (int i = 0; i < rounds; i++) {
+                CountDownLatch finished = new CountDownLatch(1);
+                StreamListener<Result> listener = new StreamListener<Result>() 
{
+                    @Override
+                    public void onNext(Result val) {
+                        // discard the placeholder prepared-statement result
+                    }
+
+                    @Override
+                    public void onError(Throwable t) {
+                        errors.incrementAndGet();
+                        finished.countDown();
+                    }
+
+                    @Override
+                    public void onCompleted() {
+                        finished.countDown();
+                    }
+                };
+                ActionCreatePreparedStatementRequest request = 
ActionCreatePreparedStatementRequest.newBuilder()
+                        .setQuery("select * from t where id = " + i).build();
+                producer.createPreparedStatement(request, callContext, 
listener);
+                Assert.assertTrue("createPreparedStatement #" + i + " did not 
finish in time",
+                        finished.await(30, TimeUnit.SECONDS));
+            }
+
+            // Guard against a false pass: if a prepare failed before reaching 
the allocation, no buffer
+            // would be leaked and the memory assertion below could not detect 
a regression.
+            Assert.assertEquals("no createPreparedStatement call should fail", 
0, errors.get());
+            // Every temporary VectorSchemaRoot must have been closed, so the 
channel's Arrow allocator
+            // is back to zero. Reverting the fix leaves `rounds` ResultMeta 
buffers allocated here.
+            Assert.assertEquals("createPreparedStatement leaked off-heap 
memory in the channel allocator",
+                    0L, channel.getAllocatedMemory());
+        } finally {
+            producer.close();
+        }
+    }
+}


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

Reply via email to