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]